brainy/src/storage/adapters/gcsStorage.ts
David Snelling 92d9420a5c fix: eliminate cloud storage write amplification and rate limiting
brain.add() was generating 26-40 immediate cloud writes per call, causing
HTTP 429 rate limit errors and high latency on GCS/S3/R2/Azure. Three-layer
fix: (1) deferred metadata writes with dirty-marking, (2) MetadataWriteBuffer
for write coalescing, (3) retry/backoff on all cloud storage adapters.
2026-01-31 09:09:36 -08:00

2119 lines
70 KiB
TypeScript

/**
* Google Cloud Storage Adapter (Native)
* Uses the native @google-cloud/storage library for optimal performance and authentication
*
* Supports multiple authentication methods:
* 1. Application Default Credentials (ADC) - Automatic in Cloud Run/GCE
* 2. Service Account Key File
* 3. Service Account Credentials Object
* 4. HMAC Keys (fallback for backward compatibility)
*/
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 { MetadataWriteBuffer } from '../../utils/metadataWriteBuffer.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'
import { InitMode } from './baseStorageAdapter.js'
// Type aliases for better readability
type HNSWNode = HNSWNoun
type Edge = HNSWVerb
// GCS client types - dynamically imported to avoid issues in browser environments
type Storage = any
type Bucket = any
type File = any
// GCS API limits
// Maximum value for maxResults parameter in GCS API calls
// Values above this cause "Invalid unsigned integer" errors
const MAX_GCS_PAGE_SIZE = 5000
/**
* Native Google Cloud Storage adapter for server environments
* Uses the @google-cloud/storage library with Application Default Credentials
*
* Authentication priority:
* 1. Application Default Credentials (if no credentials provided)
* 2. Service Account Key File (if keyFilename provided)
* 3. Service Account Credentials Object (if credentials provided)
* 4. HMAC Keys (if accessKeyId/secretAccessKey provided)
*
* 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 GcsStorage extends BaseStorage {
private storage: Storage | null = null
private bucket: Bucket | null = null
private bucketName: string
private keyFilename?: string
private credentials?: object
private accessKeyId?: string
private secretAccessKey?: 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 (_system)
// Statistics caching for better performance
protected statisticsCache: StatisticsData | null = null
// Backpressure and performance management
private pendingOperations: number = 0
private maxConcurrentOperations: number = 100
private baseBatchSize: number = 10
private currentBatchSize: number = 10
private lastMemoryCheck: number = 0
private memoryCheckInterval: number = 5000 // Check every 5 seconds
private consecutiveErrors: number = 0
private lastErrorReset: number = Date.now()
// Adaptive backpressure for automatic flow control
private backpressure = getGlobalBackpressure()
// Write buffers for bulk operations
private nounWriteBuffer: WriteBuffer<HNSWNode> | null = null
private verbWriteBuffer: WriteBuffer<Edge> | null = null
// Request coalescer for deduplication
private requestCoalescer: RequestCoalescer | null = null
// Write buffering always enabled for consistent performance
// Removes dynamic mode switching complexity - cloud storage always benefits from batching
// Multi-level cache manager for efficient data access
private nounCacheManager: CacheManager<HNSWNode>
private verbCacheManager: CacheManager<Edge>
// Module logger
private logger = createModuleLogger('GcsStorage')
// HNSW mutex locks to prevent read-modify-write races
private hnswLocks = new Map<string, Promise<void>>()
// Configuration options
private skipInitialScan: boolean = false
private skipCountsFile: boolean = false
/**
* Initialize the storage adapter
*
* @param options Configuration options for Google Cloud Storage
*
* @example Zero-config (recommended) - auto-detects Cloud Run/Lambda for fast init
* ```typescript
* const storage = new GcsStorage({ bucketName: 'my-bucket' })
* await storage.init() // <200ms in Cloud Run, blocking locally
* ```
*
* @example Force progressive mode for all environments
* ```typescript
* const storage = new GcsStorage({
* bucketName: 'my-bucket',
* initMode: 'progressive' // Always <200ms init
* })
* ```
*
* @example Force strict validation (useful for tests)
* ```typescript
* const storage = new GcsStorage({
* bucketName: 'my-bucket',
* initMode: 'strict' // Always validate during init
* })
* ```
*/
constructor(options: {
bucketName: string
// Service account authentication
keyFilename?: string
credentials?: object
// HMAC authentication (backward compatibility)
accessKeyId?: string
secretAccessKey?: string
/**
* Initialization mode for fast cold starts
*
* - `'auto'` (default): Progressive in cloud environments (Cloud Run, Lambda),
* strict locally. Zero-config optimization.
* - `'progressive'`: Always use fast init (<200ms). Bucket validation and
* count loading happen in background. First write validates bucket.
* - `'strict'`: Traditional blocking init. Validates bucket and loads counts
* before init() returns.
*
*/
initMode?: InitMode
/**
* @deprecated Use `initMode: 'progressive'` instead.
* Will be removed in a future version.
*/
skipInitialScan?: boolean
/**
* @deprecated Use `initMode: 'progressive'` instead.
* Will be removed in a future version.
*/
skipCountsFile?: boolean
// Cache and operation configuration
cacheConfig?: {
hotCacheMaxSize?: number
hotCacheEvictionThreshold?: number
warmCacheTTL?: number
}
readOnly?: boolean
}) {
super()
this.bucketName = options.bucketName
this.keyFilename = options.keyFilename
this.credentials = options.credentials
this.accessKeyId = options.accessKeyId
this.secretAccessKey = options.secretAccessKey
// Handle initMode and deprecated skip* flags
if (options.initMode) {
this.initMode = options.initMode
}
// Deprecation warnings for skip* flags
if (options.skipInitialScan) {
console.warn(
'[GcsStorage] DEPRECATION WARNING: skipInitialScan is deprecated. ' +
'Use initMode: "progressive" instead. Will be removed in a future version.'
)
this.skipInitialScan = true
}
if (options.skipCountsFile) {
console.warn(
'[GcsStorage] DEPRECATION WARNING: skipCountsFile is deprecated. ' +
'Use initMode: "progressive" instead. Will be removed in a future version.'
)
this.skipCountsFile = true
}
this.readOnly = options.readOnly || false
// 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')}/` // Noun metadata
this.verbMetadataPrefix = `${getDirectoryPath('verb', 'metadata')}/` // Verb metadata
this.systemPrefix = `${SYSTEM_DIR}/` // System data
// Initialize cache managers
this.nounCacheManager = new CacheManager<HNSWNode>(options.cacheConfig)
this.verbCacheManager = new CacheManager<Edge>(options.cacheConfig)
// Write buffering always enabled - no env var check needed
// Initialize metadata write buffer for cloud rate limit protection
this.metadataWriteBuffer = new MetadataWriteBuffer(
(path, data) => this.writeObjectToPath(path, data),
{ maxBufferSize: 200, flushIntervalMs: 200, concurrencyLimit: 10 }
)
}
/**
* Initialize the storage adapter
*
* Supports progressive initialization for fast cold starts
*
* | Mode | Init Time | When |
* |------|-----------|------|
* | `progressive` | <200ms | Cloud Run, Lambda, Azure Functions |
* | `strict` | 100-500ms+ | Local development, tests |
* | `auto` | Detected | Default - best of both |
*
* In progressive mode:
* - SDK import and bucket reference: ~50ms (unavoidable)
* - Write buffers and caches: ~10ms
* - Mark as initialized: READY
* - Background: validate bucket, load counts
*
* First write operation validates bucket existence (lazy validation).
*/
public async init(): Promise<void> {
if (this.isInitialized) {
return
}
try {
// Import Google Cloud Storage SDK only when needed (~50ms)
const { Storage } = await import('@google-cloud/storage')
// Configure the GCS client based on available credentials
const clientConfig: any = {}
// Priority 1: Service Account Key File
if (this.keyFilename) {
clientConfig.keyFilename = this.keyFilename
prodLog.info('🔐 GCS: Using Service Account Key File')
}
// Priority 2: Service Account Credentials Object
else if (this.credentials) {
clientConfig.credentials = this.credentials
prodLog.info('🔐 GCS: Using Service Account Credentials')
}
// Priority 3: HMAC Keys (S3 compatibility)
else if (this.accessKeyId && this.secretAccessKey) {
clientConfig.credentials = {
client_email: 'hmac-user@example.com',
private_key: this.secretAccessKey
}
prodLog.warn('⚠️ GCS: Using HMAC keys (consider migrating to ADC)')
}
// Priority 4: Application Default Credentials (default)
else {
// No credentials needed - ADC will be used automatically
prodLog.info('🔐 GCS: Using Application Default Credentials (ADC)')
}
// Create the GCS client (no network calls)
this.storage = new Storage(clientConfig)
// Get reference to the bucket (no network calls)
this.bucket = this.storage.bucket(this.bucketName)
// Determine initialization mode
const effectiveMode = this.resolveInitMode()
const isCloud = this.detectCloudEnvironment()
prodLog.info(`🚀 GCS init mode: ${effectiveMode} (detected cloud: ${isCloud})`)
// Initialize write buffers for high-volume mode
const storageId = `gcs-${this.bucketName}`
this.nounWriteBuffer = getWriteBuffer<HNSWNode>(
`${storageId}-nouns`,
'noun',
async (items) => {
await this.flushNounBuffer(items)
}
)
this.verbWriteBuffer = getWriteBuffer<Edge>(
`${storageId}-verbs`,
'verb',
async (items) => {
await this.flushVerbBuffer(items)
}
)
// Initialize request coalescer for deduplication
this.requestCoalescer = getCoalescer(
storageId,
async (batch) => {
// Process coalesced operations (placeholder for future optimization)
this.logger.trace(`Processing coalesced batch: ${batch.length} items`)
}
)
// CRITICAL FIX: Clear any stale cache entries from previous runs
// This prevents cache poisoning from causing silent failures on container restart
prodLog.info('🧹 Clearing cache from previous run to prevent cache poisoning')
this.nounCacheManager.clear()
this.verbCacheManager.clear()
prodLog.info('✅ Cache cleared - starting fresh')
// Progressive vs Strict initialization
if (effectiveMode === 'progressive') {
// PROGRESSIVE MODE: Fast init, background validation
// Mark as initialized immediately - ready to accept operations
// Bucket validation happens lazily on first write
// Count loading happens in background
prodLog.info(`✅ GCS progressive init complete: ${this.bucketName} (validation deferred)`)
// Initialize GraphAdjacencyIndex and type statistics
await super.init()
// Schedule background tasks (non-blocking)
this.scheduleBackgroundInit()
} else {
// STRICT MODE: Traditional blocking initialization
// Verify bucket exists and is accessible (blocking)
const [exists] = await this.bucket.exists()
if (!exists) {
throw new Error(`Bucket ${this.bucketName} does not exist or is not accessible`)
}
this.bucketValidated = true
prodLog.info(`✅ Connected to GCS bucket: ${this.bucketName}`)
// Initialize counts from storage (blocking)
await this.initializeCounts()
this.countsLoaded = true
// Initialize GraphAdjacencyIndex and type statistics
await super.init()
// Mark background tasks as complete (nothing to do in background)
this.backgroundTasksComplete = true
}
} catch (error) {
this.logger.error('Failed to initialize GCS storage:', error)
throw new Error(`Failed to initialize GCS storage: ${error}`)
}
}
// =============================================
// Progressive Initialization
// =============================================
/**
* Run background initialization tasks for GCS.
*
* Called in progressive mode after init() returns. Performs:
* 1. Bucket validation (in background)
* 2. Count loading from storage (in background)
*
* These tasks don't block the main thread, allowing fast cold starts.
*
* @protected
* @override
*/
protected async runBackgroundInit(): Promise<void> {
const startTime = Date.now()
prodLog.info('[GCS Background] Starting background initialization...')
// Run validation and count loading in parallel
const validationPromise = this.validateBucketInBackground()
const countsPromise = this.loadCountsInBackground()
// Wait for both to complete (but we're already initialized)
await Promise.all([validationPromise, countsPromise])
const elapsed = Date.now() - startTime
prodLog.info(`[GCS Background] Background init complete in ${elapsed}ms`)
}
/**
* Validate bucket existence in background.
*
* Stores result in bucketValidated/bucketValidationError for lazy use.
*
* @private
*/
private async validateBucketInBackground(): Promise<void> {
try {
const [exists] = await this.bucket!.exists()
if (exists) {
this.bucketValidated = true
prodLog.info(`[GCS Background] Bucket validated: ${this.bucketName}`)
} else {
this.bucketValidationError = new Error(
`Bucket ${this.bucketName} does not exist or is not accessible`
)
prodLog.warn(`[GCS Background] Bucket validation failed: ${this.bucketName}`)
}
} catch (error: any) {
this.bucketValidationError = new Error(
`Bucket validation failed: ${error.message || error}`
)
prodLog.error(`[GCS Background] Bucket validation error:`, error)
}
}
/**
* Load counts from storage in background.
*
* Uses the existing initializeCounts() logic but in background.
*
* @private
*/
private async loadCountsInBackground(): Promise<void> {
try {
await this.initializeCounts()
this.countsLoaded = true
prodLog.info(`[GCS Background] Counts loaded: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs`)
} catch (error: any) {
// Non-fatal in progressive mode - counts start at 0
prodLog.warn(`[GCS Background] Failed to load counts (starting at 0): ${error.message}`)
this.countsLoaded = true // Mark as loaded even on error (0 is valid)
}
}
/**
* Ensure bucket is validated before write operations.
*
* In progressive mode, bucket validation is deferred until the first
* write operation. This method validates the bucket and caches the
* result (or error) for subsequent calls.
*
* @throws Error if bucket does not exist or is not accessible
* @protected
* @override
*/
protected async ensureValidatedForWrite(): Promise<void> {
// If already validated, nothing to do
if (this.bucketValidated) {
return
}
// If we have a cached validation error from background init, throw it
if (this.bucketValidationError) {
throw this.bucketValidationError
}
// Perform synchronous validation (first write in progressive mode)
try {
prodLog.info(`[GCS] Lazy bucket validation on first write: ${this.bucketName}`)
const [exists] = await this.bucket!.exists()
if (!exists) {
const error = new Error(`Bucket ${this.bucketName} does not exist or is not accessible`)
this.bucketValidationError = error
throw error
}
this.bucketValidated = true
prodLog.info(`[GCS] Bucket validated successfully: ${this.bucketName}`)
} catch (error: any) {
// Cache the error for fast-fail on subsequent writes
if (!this.bucketValidationError) {
this.bucketValidationError = error
}
throw error
}
}
/**
* Get the GCS object key for a noun using UUID-based sharding
*
* Uses first 2 hex characters of UUID for consistent sharding.
* Path format: entities/nouns/vectors/{shardId}/{uuid}.json
*
* @example
* getNounKey('ab123456-1234-5678-9abc-def012345678')
* // returns 'entities/nouns/vectors/ab/ab123456-1234-5678-9abc-def012345678.json'
*/
private getNounKey(id: string): string {
const shardId = getShardIdFromUuid(id)
return `${this.nounPrefix}${shardId}/${id}.json`
}
/**
* Get the GCS object key for a verb using UUID-based sharding
*
* Uses first 2 hex characters of UUID for consistent sharding.
* Path format: entities/verbs/vectors/{shardId}/{uuid}.json
*
* @example
* getVerbKey('cd987654-4321-8765-cba9-fed543210987')
* // returns 'entities/verbs/vectors/cd/cd987654-4321-8765-cba9-fed543210987.json'
*/
private getVerbKey(id: string): string {
const shardId = getShardIdFromUuid(id)
return `${this.verbPrefix}${shardId}/${id}.json`
}
/**
* Override base class method to detect GCS-specific throttling errors
*/
protected isThrottlingError(error: any): boolean {
// First check base class detection
if (super.isThrottlingError(error)) {
return true
}
// GCS-specific throttling detection
const statusCode = error.code
const message = error.message?.toLowerCase() || ''
return (
statusCode === 429 || // Too Many Requests
statusCode === 503 || // Service Unavailable
statusCode === 'RATE_LIMIT_EXCEEDED' ||
message.includes('quota') ||
message.includes('rate limit') ||
message.includes('too many requests')
)
}
/**
* Override base class to enable smart batching for cloud storage
*
* GCS is cloud storage with network latency (~50ms per write).
* Smart batching reduces writes from 1000 ops → 100 batches.
*
* @returns true (GCS is cloud storage)
*/
protected isCloudStorage(): boolean {
return true // GCS benefits from batching
}
/**
* Apply backpressure before starting an operation
* @returns Request ID for tracking
*/
private async applyBackpressure(): Promise<string> {
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
* @param success Whether the operation succeeded
* @param requestId Request ID from applyBackpressure()
*/
private releaseBackpressure(success: boolean = true, requestId?: string): void {
this.pendingOperations = Math.max(0, this.pendingOperations - 1)
if (requestId) {
this.backpressure.releasePermission(requestId, success)
}
}
// Removed checkVolumeMode() - write buffering always enabled for cloud storage
/**
* Flush noun buffer to GCS
*/
private async flushNounBuffer(items: Map<string, HNSWNode>): Promise<void> {
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 GCS
*/
private async flushVerbBuffer(items: Map<string, Edge>): Promise<void> {
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)
}
// Removed saveNoun_internal - now inherit from BaseStorage's type-first implementation
/**
* Save a node to storage
* Always uses write buffer for consistent performance
*/
protected async saveNode(node: HNSWNode): Promise<void> {
await this.ensureInitialized()
// Always use write buffer - cloud storage benefits from batching
if (this.nounWriteBuffer) {
this.logger.trace(`📝 BUFFERING: Adding noun ${node.id} to write buffer`)
// Populate cache BEFORE buffering for read-after-write consistency
if (node.vector && Array.isArray(node.vector) && node.vector.length > 0) {
this.nounCacheManager.set(node.id, node)
}
await this.nounWriteBuffer.add(node.id, node)
return
}
// Fallback to direct write if buffer not initialized (shouldn't happen after init)
await this.saveNodeDirect(node)
}
/**
* Save a node directly to GCS (bypass buffer)
*
* Performs lazy bucket validation on first write in progressive mode.
*/
private async saveNodeDirect(node: HNSWNode): Promise<void> {
// Lazy bucket validation for progressive init
await this.ensureValidatedForWrite()
// Apply backpressure before starting operation
const requestId = await this.applyBackpressure()
try {
this.logger.trace(`Saving node ${node.id}`)
// Convert connections Map to a serializable format
// CRITICAL: Only save lightweight vector data (no metadata)
// Metadata is saved separately via saveNounMetadata() (2-file system)
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
// NO metadata field - saved separately for scalability
}
// Get the GCS key with UUID-based sharding
const key = this.getNounKey(node.id)
// Save to GCS
const file = this.bucket!.file(key)
await file.save(JSON.stringify(serializableNode, null, 2), {
contentType: 'application/json',
resumable: false // For small objects, non-resumable is faster
})
// CRITICAL FIX: Only cache nodes with non-empty vectors
// This prevents cache pollution from HNSW's lazy-loading nodes (vector: [])
if (node.vector && Array.isArray(node.vector) && node.vector.length > 0) {
this.nounCacheManager.set(node.id, node)
}
// Note: Empty vectors are intentional during HNSW lazy mode - not logged
// Count tracking happens in baseStorage.saveNounMetadata_internal
// This fixes the race condition where metadata didn't exist yet
this.logger.trace(`Node ${node.id} saved successfully`)
this.releaseBackpressure(true, requestId)
} catch (error: any) {
this.releaseBackpressure(false, requestId)
// Handle throttling
if (this.isThrottlingError(error)) {
await this.handleThrottling(error)
throw error // Re-throw for retry at higher level
}
this.logger.error(`Failed to save node ${node.id}:`, error)
throw new Error(`Failed to save node ${node.id}: ${error}`)
}
}
// Removed getNoun_internal - now inherit from BaseStorage's type-first implementation
/**
* Get a node from storage
*/
protected async getNode(id: string): Promise<HNSWNode | null> {
await this.ensureInitialized()
// Check cache first
const cached: HNSWNode | null = await this.nounCacheManager.get(id)
// Validate cached object before returning
if (cached !== undefined && cached !== null) {
// Validate cached object has required fields (including non-empty vector!)
if (!cached.id || !cached.vector || !Array.isArray(cached.vector) || cached.vector.length === 0) {
// Invalid cache detected - log and auto-recover
prodLog.warn(`[GCS] Invalid cached object for ${id.substring(0, 8)} (${
!cached.id ? 'missing id' :
!cached.vector ? 'missing vector' :
!Array.isArray(cached.vector) ? 'vector not array' :
'empty vector'
}) - removing from cache and reloading`)
this.nounCacheManager.delete(id)
// Fall through to load from GCS
} else {
// Valid cache hit
this.logger.trace(`Cache hit for noun ${id}`)
return cached
}
} else if (cached === null) {
prodLog.warn(`[GCS] Cache contains null for ${id.substring(0, 8)} - reloading from storage`)
}
// Apply backpressure
const requestId = await this.applyBackpressure()
try {
this.logger.trace(`Getting node ${id}`)
// Get the GCS key with UUID-based sharding
const key = this.getNounKey(id)
// Download from GCS
const file = this.bucket!.file(key)
const [contents] = await file.download()
// Parse JSON
const data = JSON.parse(contents.toString())
// Convert serialized connections back to Map<number, Set<string>>
const connections = new Map<number, Set<string>>()
for (const [level, nounIds] of Object.entries(data.connections || {})) {
connections.set(Number(level), new Set(nounIds as string[]))
}
// CRITICAL: Only return lightweight vector data (no metadata)
// Metadata is retrieved separately via getNounMetadata() (2-file system)
const node: HNSWNode = {
id: data.id,
vector: data.vector,
connections,
level: data.level || 0
// NO metadata field - retrieved separately for scalability
}
// CRITICAL FIX: Only cache valid nodes with non-empty vectors (never cache null or empty)
if (node && node.id && node.vector && Array.isArray(node.vector) && node.vector.length > 0) {
this.nounCacheManager.set(id, node)
} else {
prodLog.warn(`[GCS] Not caching invalid node ${id.substring(0, 8)} (missing id/vector or empty vector)`)
}
this.logger.trace(`Successfully retrieved node ${id}`)
this.releaseBackpressure(true, requestId)
return node
} catch (error: any) {
this.releaseBackpressure(false, requestId)
// DIAGNOSTIC LOGGING: Log EVERY error before any conditional checks
const key = this.getNounKey(id)
prodLog.error(`[getNode] ❌ EXCEPTION CAUGHT:`)
prodLog.error(`[getNode] UUID: ${id}`)
prodLog.error(`[getNode] Path: ${key}`)
prodLog.error(`[getNode] Bucket: ${this.bucketName}`)
prodLog.error(`[getNode] Error type: ${error?.constructor?.name || typeof error}`)
prodLog.error(`[getNode] Error code: ${JSON.stringify(error?.code)}`)
prodLog.error(`[getNode] Error message: ${error?.message || String(error)}`)
prodLog.error(`[getNode] Error object:`, JSON.stringify(error, null, 2))
// Check if this is a "not found" error
if (error.code === 404) {
prodLog.warn(`[getNode] Identified as 404 error - returning null WITHOUT caching`)
// CRITICAL FIX: Do NOT cache null values
return null
}
// Handle throttling
if (this.isThrottlingError(error)) {
prodLog.warn(`[getNode] Identified as throttling error - rethrowing`)
await this.handleThrottling(error)
throw error
}
// All other errors should throw, not return null
prodLog.error(`[getNode] Unhandled error - rethrowing`)
this.logger.error(`Failed to get node ${id}:`, error)
throw BrainyError.fromError(error, `getNoun(${id})`)
}
}
// Removed deleteNoun_internal - now inherit from BaseStorage's type-first implementation
/**
* Write an object to a specific path in GCS
* Primitive operation required by base class
*
* Performs lazy bucket validation on first write in progressive mode.
* @protected
*/
protected async writeObjectToPath(path: string, data: any): Promise<void> {
await this.ensureInitialized()
// Lazy bucket validation for progressive init
await this.ensureValidatedForWrite()
const MAX_RETRIES = 5
let lastError: any
for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) {
try {
this.logger.trace(`Writing object to path: ${path}`)
const file = this.bucket!.file(path)
await file.save(JSON.stringify(data, null, 2), {
contentType: 'application/json',
resumable: false
})
this.logger.trace(`Object written successfully to ${path}`)
if (attempt > 0) {
this.clearThrottlingState()
}
return
} catch (error: any) {
lastError = error
if (this.isThrottlingError(error) && attempt < MAX_RETRIES) {
const baseDelay = Math.min(100 * Math.pow(2, attempt), 5000)
const jitter = baseDelay * 0.2 * (Math.random() - 0.5)
const delay = Math.round(baseDelay + jitter)
this.logger.warn(
`Throttled writing ${path} (attempt ${attempt + 1}/${MAX_RETRIES + 1}), retrying in ${delay}ms`
)
await this.handleThrottling(error)
await new Promise(resolve => setTimeout(resolve, delay))
continue
}
break
}
}
this.logger.error(`Failed to write object to ${path}:`, lastError)
throw new Error(`Failed to write object to ${path}: ${lastError}`)
}
/**
* Read an object from a specific path in GCS
* Primitive operation required by base class
* @protected
*/
protected async readObjectFromPath(path: string): Promise<any | null> {
await this.ensureInitialized()
try {
this.logger.trace(`Reading object from path: ${path}`)
const file = this.bucket!.file(path)
const [contents] = await file.download()
const data = JSON.parse(contents.toString())
this.logger.trace(`Object read successfully from ${path}`)
return data
} catch (error: any) {
// Check if this is a "not found" error
if (error.code === 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})`)
}
}
/**
* Batch read multiple objects from GCS
*
* **Performance**: GCS-optimized parallel downloads
* - Uses Promise.all() for concurrent requests
* - Respects GCS rate limits (100 concurrent by default)
* - Chunks large batches to prevent memory issues
*
* **GCS Specifics**:
* - No true "batch API" - uses parallel GetObject operations
* - Optimal concurrency: 50-100 concurrent downloads
* - Each download is a separate HTTPS request
*
* @param paths Array of GCS object paths to read
* @returns Map of path → data (only successful reads included)
*
* @public - Called by baseStorage.readBatchFromAdapter()
*/
public async readBatch(paths: string[]): Promise<Map<string, any>> {
await this.ensureInitialized()
const results = new Map<string, any>()
if (paths.length === 0) return results
// Get batch configuration for optimal GCS performance
const batchConfig = this.getBatchConfig()
const chunkSize = batchConfig.maxConcurrent || 100
this.logger.debug(`[GCS Batch] Reading ${paths.length} objects in chunks of ${chunkSize}`)
// Process in chunks to respect rate limits and prevent memory issues
for (let i = 0; i < paths.length; i += chunkSize) {
const chunk = paths.slice(i, i + chunkSize)
this.logger.trace(`[GCS Batch] Processing chunk ${Math.floor(i/chunkSize) + 1}/${Math.ceil(paths.length/chunkSize)}`)
// Parallel download for this chunk
const chunkResults = await Promise.allSettled(
chunk.map(async (path) => {
try {
const file = this.bucket!.file(path)
const [contents] = await file.download()
const data = JSON.parse(contents.toString())
return { path, data, success: true }
} catch (error: any) {
// Silently skip 404s (expected for missing entities)
if (error.code === 404) {
return { path, data: null, success: false }
}
// Log other errors but don't fail the batch
this.logger.warn(`[GCS Batch] Failed to read ${path}: ${error.message}`)
return { path, data: null, success: false }
}
})
)
// Collect successful results
for (const result of chunkResults) {
if (result.status === 'fulfilled' && result.value.success && result.value.data !== null) {
results.set(result.value.path, result.value.data)
}
}
}
this.logger.debug(`[GCS Batch] Successfully read ${results.size}/${paths.length} objects`)
return results
}
/**
* Get GCS-specific batch configuration
*
* GCS performs well with high concurrency due to HTTP/2 multiplexing
*
* @public - Overrides BaseStorage.getBatchConfig()
*/
public getBatchConfig(): StorageBatchConfig {
return {
maxBatchSize: 1000, // GCS can handle large batches
batchDelayMs: 0, // No rate limiting needed (HTTP/2 handles it)
maxConcurrent: 100, // Optimal for GCS (tested up to 200)
supportsParallelWrites: true,
rateLimit: {
operationsPerSecond: 1000, // GCS is fast
burstCapacity: 5000
}
}
}
/**
* Delete an object from a specific path in GCS
* Primitive operation required by base class
*
* Performs lazy bucket validation on first delete in progressive mode.
* @protected
*/
protected async deleteObjectFromPath(path: string): Promise<void> {
await this.ensureInitialized()
// Lazy bucket validation for progressive init
await this.ensureValidatedForWrite()
try {
this.logger.trace(`Deleting object at path: ${path}`)
const file = this.bucket!.file(path)
await file.delete()
this.logger.trace(`Object deleted successfully from ${path}`)
} catch (error: any) {
// If already deleted (404), treat as success
if (error.code === 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 GCS
* Primitive operation required by base class
* @protected
*/
protected async listObjectsUnderPath(prefix: string): Promise<string[]> {
await this.ensureInitialized()
try {
this.logger.trace(`Listing objects under prefix: ${prefix}`)
const [files] = await this.bucket!.getFiles({ prefix })
const paths = files.map((file: any) => file.name).filter((name: string) => name && name.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}`)
}
}
// Removed saveVerb_internal - now inherit from BaseStorage's type-first implementation
/**
* Save an edge to storage
* Always uses write buffer for consistent performance
*/
protected async saveEdge(edge: Edge): Promise<void> {
await this.ensureInitialized()
// Always use write buffer - cloud storage benefits from batching
if (this.verbWriteBuffer) {
this.logger.trace(`📝 BUFFERING: Adding verb ${edge.id} to write buffer`)
// Populate cache BEFORE buffering for read-after-write consistency
this.verbCacheManager.set(edge.id, edge)
await this.verbWriteBuffer.add(edge.id, edge)
return
}
// Fallback to direct write if buffer not initialized (shouldn't happen after init)
await this.saveEdgeDirect(edge)
}
/**
* Save an edge directly to GCS (bypass buffer)
*
* Performs lazy bucket validation on first write in progressive mode.
*/
private async saveEdgeDirect(edge: Edge): Promise<void> {
// Lazy bucket validation for progressive init
await this.ensureValidatedForWrite()
const requestId = await this.applyBackpressure()
try {
this.logger.trace(`Saving edge ${edge.id}`)
// Convert connections Map to serializable format
// ARCHITECTURAL FIX: 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
verb: edge.verb,
sourceId: edge.sourceId,
targetId: edge.targetId,
// User metadata (if any) - saved separately for scalability
// metadata field is saved separately via saveVerbMetadata()
}
// Get the GCS key with UUID-based sharding
const key = this.getVerbKey(edge.id)
// Save to GCS
const file = this.bucket!.file(key)
await file.save(JSON.stringify(serializableEdge, null, 2), {
contentType: 'application/json',
resumable: false
})
// Update cache
this.verbCacheManager.set(edge.id, edge)
// Count tracking happens in baseStorage.saveVerbMetadata_internal
// This fixes the race condition where metadata didn't exist yet
this.logger.trace(`Edge ${edge.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 edge ${edge.id}:`, error)
throw new Error(`Failed to save edge ${edge.id}: ${error}`)
}
}
// Removed getVerb_internal - now inherit from BaseStorage's type-first implementation
/**
* Get an edge from storage
*/
protected async getEdge(id: string): Promise<Edge | null> {
await this.ensureInitialized()
// Check cache first
const cached = this.verbCacheManager.get(id)
if (cached) {
this.logger.trace(`Cache hit for verb ${id}`)
return cached
}
const requestId = await this.applyBackpressure()
try {
this.logger.trace(`Getting edge ${id}`)
// Get the GCS key with UUID-based sharding
const key = this.getVerbKey(id)
// Download from GCS
const file = this.bucket!.file(key)
const [contents] = await file.download()
// Parse JSON
const data = JSON.parse(contents.toString())
// Convert serialized connections back to Map
const connections = new Map<number, Set<string>>()
for (const [level, verbIds] of Object.entries(data.connections || {})) {
connections.set(Number(level), new Set(verbIds as string[]))
}
// 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
// User metadata retrieved separately via getVerbMetadata()
}
// Update cache
this.verbCacheManager.set(id, edge)
this.logger.trace(`Successfully retrieved edge ${id}`)
this.releaseBackpressure(true, requestId)
return edge
} catch (error: any) {
this.releaseBackpressure(false, requestId)
// Check if this is a "not found" error
if (error.code === 404) {
this.logger.trace(`Edge not found: ${id}`)
return null
}
if (this.isThrottlingError(error)) {
await this.handleThrottling(error)
throw error
}
this.logger.error(`Failed to get edge ${id}:`, error)
throw BrainyError.fromError(error, `getVerb(${id})`)
}
}
// Removed deleteVerb_internal - now inherit from BaseStorage's type-first implementation
// Removed pagination overrides - use BaseStorage's type-first implementation
// - getNounsWithPagination, getNodesWithPagination, getVerbsWithPagination
// - getNouns, getVerbs (public wrappers)
// Removed 4 query *_internal methods - now inherit from BaseStorage's type-first implementation
// (getNounsByNounType_internal, getVerbsBySource_internal, getVerbsByTarget_internal, getVerbsByType_internal)
/**
* Batch fetch metadata for multiple noun IDs (efficient for large queries)
* Uses smaller batches to prevent GCS socket exhaustion
* @param ids Array of noun IDs to fetch metadata for
* @returns Map of ID to metadata
*/
public async getMetadataBatch(ids: string[]): Promise<Map<string, any>> {
await this.ensureInitialized()
const results = new Map<string, any>()
const batchSize = 10 // Smaller batches for metadata to prevent socket exhaustion
// Process in smaller batches
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: any) {
// Handle GCS-specific errors
if (this.isThrottlingError(error)) {
await this.handleThrottling(error)
}
this.logger.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 to prevent overwhelming GCS
await new Promise(resolve => setImmediate(resolve))
}
return results
}
/**
* Clear all data from storage
*/
public async clear(): Promise<void> {
await this.ensureInitialized()
try {
this.logger.info('🧹 Clearing all data from GCS bucket...')
// Helper function to delete all objects with a given prefix
const deleteObjectsWithPrefix = async (prefix: string): Promise<void> => {
const [files] = await this.bucket!.getFiles({ prefix })
if (!files || files.length === 0) {
return
}
// Delete each file
for (const file of files) {
await file.delete()
}
}
// Clear ALL data using correct paths
// Delete entire branches/ directory (includes ALL entities, ALL types, ALL VFS data, ALL forks)
await deleteObjectsWithPrefix('branches/')
// Delete COW version control data
await deleteObjectsWithPrefix('_cow/')
// Delete system metadata
await deleteObjectsWithPrefix('_system/')
// Reset COW managers (but don't disable COW - it's always enabled)
// COW will re-initialize automatically on next use
this.refManager = undefined
this.blobStorage = undefined
this.commitLog = undefined
// Clear caches
this.nounCacheManager.clear()
this.verbCacheManager.clear()
// Reset counts
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts.clear()
this.verbCounts.clear()
this.logger.info('✅ All data cleared from GCS')
} catch (error) {
this.logger.error('Failed to clear GCS storage:', error)
throw new Error(`Failed to clear GCS storage: ${error}`)
}
}
/**
* Get storage status
*/
public async getStorageStatus(): Promise<{
type: string
used: number
quota: number | null
details?: Record<string, any>
}> {
await this.ensureInitialized()
try {
// Get bucket metadata
const [metadata] = await this.bucket!.getMetadata()
return {
type: 'gcs', // Consistent with new naming (native SDK is just 'gcs')
used: 0, // GCS doesn't provide usage info easily
quota: null, // No quota in GCS
details: {
bucket: this.bucketName,
location: metadata.location,
storageClass: metadata.storageClass,
created: metadata.timeCreated,
sdk: 'native' // Indicate we're using native SDK
}
}
} catch (error) {
this.logger.error('Failed to get storage status:', error)
return {
type: 'gcs', // Consistent with new naming
used: 0,
quota: null
}
}
}
/**
* Check if COW has been explicitly disabled via clear()
* Fixes bug where clear() doesn't persist across instance restarts
* @returns true if marker object exists, false otherwise
* @protected
*/
/**
* Removed checkClearMarker() and createClearMarker() methods
* COW is now always enabled - marker files are no longer used
*/
/**
* Save statistics data to storage
*/
protected async saveStatisticsData(statistics: StatisticsData): Promise<void> {
await this.ensureInitialized()
try {
const key = `${this.systemPrefix}${STATISTICS_KEY}.json`
this.logger.trace(`Saving statistics to ${key}`)
const file = this.bucket!.file(key)
await file.save(JSON.stringify(statistics, null, 2), {
contentType: 'application/json',
resumable: false
})
this.logger.trace('Statistics saved successfully')
} catch (error) {
this.logger.error('Failed to save statistics:', error)
throw new Error(`Failed to save statistics: ${error}`)
}
}
/**
* Get statistics data from storage
*/
protected async getStatisticsData(): Promise<StatisticsData | null> {
await this.ensureInitialized()
try {
const key = `${this.systemPrefix}${STATISTICS_KEY}.json`
this.logger.trace(`Getting statistics from ${key}`)
const file = this.bucket!.file(key)
const [contents] = await file.download()
const statistics = JSON.parse(contents.toString())
this.logger.trace('Statistics retrieved successfully')
// CRITICAL FIX: Populate totalNodes and totalEdges from in-memory counts
// HNSW rebuild depends on these fields to determine entity count
return {
...statistics,
totalNodes: this.totalNounCount,
totalEdges: this.totalVerbCount,
lastUpdated: new Date().toISOString()
}
} catch (error: any) {
if (error.code === 404) {
// CRITICAL FIX: Statistics file doesn't exist yet (first restart)
// Return minimal stats with counts instead of null
// This prevents HNSW from seeing entityCount=0 during index rebuild
this.logger.trace('Statistics file not found - returning minimal stats with counts')
return {
nounCount: {},
verbCount: {},
metadataCount: {},
hnswIndexSize: 0,
totalNodes: this.totalNounCount,
totalEdges: this.totalVerbCount,
totalMetadata: 0,
lastUpdated: new Date().toISOString()
}
}
this.logger.error('Failed to get statistics:', error)
return null
}
}
/**
* Initialize counts from storage
*/
protected async initializeCounts(): Promise<void> {
// Skip counts file entirely if configured
if (this.skipCountsFile) {
prodLog.info('📊 Skipping counts file (skipCountsFile: true)')
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts = new Map()
this.verbCounts = new Map()
return
}
const key = `${this.systemPrefix}counts.json`
try {
const file = this.bucket!.file(key)
const [contents] = await file.download()
const counts = JSON.parse(contents.toString())
this.totalNounCount = counts.totalNounCount || 0
this.totalVerbCount = counts.totalVerbCount || 0
this.entityCounts = new Map(Object.entries(counts.entityCounts || {})) as Map<string, number>
this.verbCounts = new Map(Object.entries(counts.verbCounts || {})) as Map<string, number>
prodLog.info(`📊 Loaded counts from storage: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs`)
} catch (error: any) {
if (error.code === 404) {
// No counts file yet
if (this.skipInitialScan) {
prodLog.info('📊 No counts file found - starting with zero counts (skipInitialScan: true)')
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts = new Map()
this.verbCounts = new Map()
} else {
// Initialize from scan (first-time setup or counts not persisted)
prodLog.info('📊 No counts file found - scanning bucket to initialize counts')
await this.initializeCountsFromScan()
}
} else {
// CRITICAL FIX: Don't silently fail on network/permission errors
this.logger.error('❌ CRITICAL: Failed to load counts from GCS:', error)
prodLog.error(`❌ Error loading ${key}: ${error.message}`)
if (this.skipInitialScan) {
prodLog.warn('⚠️ Starting with zero counts due to error (skipInitialScan: true)')
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts = new Map()
this.verbCounts = new Map()
} else {
// Try to recover by scanning the bucket
prodLog.warn('⚠️ Attempting recovery by scanning GCS bucket...')
await this.initializeCountsFromScan()
}
}
}
}
/**
* Initialize counts from storage scan (expensive - only for first-time init)
* Includes timeout handling to prevent Cloud Run startup failures
*/
private async initializeCountsFromScan(): Promise<void> {
const SCAN_TIMEOUT_MS = 120000 // 2 minutes timeout
try {
prodLog.info('📊 Scanning GCS bucket to initialize counts...')
prodLog.info(`🔍 Noun prefix: ${this.nounPrefix}`)
prodLog.info(`🔍 Verb prefix: ${this.verbPrefix}`)
prodLog.info(`⏱️ Timeout: ${SCAN_TIMEOUT_MS / 1000}s (configure skipInitialScan to avoid this)`)
// Create timeout promise
const timeoutPromise = new Promise<never>((_, reject) => {
setTimeout(() => {
reject(new Error(`Bucket scan timeout after ${SCAN_TIMEOUT_MS / 1000}s`))
}, SCAN_TIMEOUT_MS)
})
// Count nouns with timeout
const nounScanPromise = this.bucket!.getFiles({ prefix: this.nounPrefix })
const [nounFiles] = await Promise.race([nounScanPromise, timeoutPromise]) as any
prodLog.info(`🔍 Found ${nounFiles?.length || 0} total files under noun prefix`)
const jsonNounFiles = nounFiles?.filter((f: any) => f.name?.endsWith('.json')) || []
this.totalNounCount = jsonNounFiles.length
if (jsonNounFiles.length > 0 && jsonNounFiles.length <= 5) {
prodLog.info(`📄 Sample noun files: ${jsonNounFiles.slice(0, 5).map((f: any) => f.name).join(', ')}`)
}
// Count verbs with timeout
const verbScanPromise = this.bucket!.getFiles({ prefix: this.verbPrefix })
const [verbFiles] = await Promise.race([verbScanPromise, timeoutPromise]) as any
prodLog.info(`🔍 Found ${verbFiles?.length || 0} total files under verb prefix`)
const jsonVerbFiles = verbFiles?.filter((f: any) => f.name?.endsWith('.json')) || []
this.totalVerbCount = jsonVerbFiles.length
if (jsonVerbFiles.length > 0 && jsonVerbFiles.length <= 5) {
prodLog.info(`📄 Sample verb files: ${jsonVerbFiles.slice(0, 5).map((f: any) => f.name).join(', ')}`)
}
// Save initial counts
if (this.totalNounCount > 0 || this.totalVerbCount > 0) {
await this.persistCounts()
prodLog.info(`✅ Initialized counts from scan: ${this.totalNounCount} nouns, ${this.totalVerbCount} verbs`)
} else {
prodLog.warn(`⚠️ No entities found during bucket scan. Check that entities exist and prefixes are correct.`)
}
} catch (error: any) {
// Handle timeout specifically
if (error.message?.includes('Bucket scan timeout')) {
prodLog.error(`❌ TIMEOUT: Bucket scan exceeded ${SCAN_TIMEOUT_MS / 1000}s limit`)
prodLog.error(` This typically happens with large buckets in Cloud Run deployments.`)
prodLog.error(` Solutions:`)
prodLog.error(` 1. Increase Cloud Run timeout: timeoutSeconds: 600`)
prodLog.error(` 2. Use skipInitialScan: true in gcsNativeStorage config`)
prodLog.error(` 3. Pre-create counts file before deployment`)
prodLog.warn(`⚠️ Starting with zero counts due to timeout`)
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts = new Map()
this.verbCounts = new Map()
return
}
// CRITICAL FIX: Don't silently fail - this prevents data loss scenarios
this.logger.error('❌ CRITICAL: Failed to initialize counts from GCS bucket scan:', error)
prodLog.error(` Error: ${error.message || String(error)}`)
prodLog.warn(`⚠️ Starting with zero counts due to error`)
this.totalNounCount = 0
this.totalVerbCount = 0
this.entityCounts = new Map()
this.verbCounts = new Map()
}
}
/**
* Persist counts to storage
*/
protected async persistCounts(): Promise<void> {
// Skip if skipCountsFile is enabled
if (this.skipCountsFile) {
return
}
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()
}
const file = this.bucket!.file(key)
await file.save(JSON.stringify(counts, null, 2), {
contentType: 'application/json',
resumable: false
})
} catch (error) {
this.logger.error('Error persisting counts:', error)
}
}
// HNSW Index Persistence
/**
* Get a noun's vector for HNSW rebuild
* Uses BaseStorage's getNoun (type-first paths)
*/
public async getNounVector(id: string): Promise<number[] | null> {
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: Uses mutex locking to prevent read-modify-write races
*/
public async saveHNSWData(nounId: string, hnswData: {
level: number
connections: Record<string, string[]>
}): Promise<void> {
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 FileSystem/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<void>(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<number, Set<string>>()
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 via writeObjectToBranch)
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 getHNSWData(nounId: string): Promise<{
level: number
connections: Record<string, string[]>
} | null> {
const noun = await this.getNoun(nounId)
if (!noun) {
return null
}
// Convert connections from Map to Record format
const connectionsRecord: Record<string, string[]> = {}
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)
* Storage path: system/hnsw-system.json
*
* CRITICAL FIX: Optimistic locking with generation numbers to prevent race conditions
*/
public async saveHNSWSystem(systemData: {
entryPointId: string | null
maxLevel: number
}): Promise<void> {
await this.ensureInitialized()
const key = `${this.systemPrefix}hnsw-system.json`
const file = this.bucket!.file(key)
const maxRetries = 5
for (let attempt = 0; attempt < maxRetries; attempt++) {
try {
// Get current generation
let currentGeneration: string | undefined
try {
const [metadata] = await file.getMetadata()
currentGeneration = metadata.generation?.toString()
} catch (error: any) {
// File doesn't exist yet
if (error.code !== 404) {
throw error
}
}
// ATOMIC WRITE: Use generation precondition
await file.save(JSON.stringify(systemData, null, 2), {
contentType: 'application/json',
resumable: false,
preconditionOpts: currentGeneration
? { ifGenerationMatch: currentGeneration }
: { ifGenerationMatch: '0' }
})
// Success!
return
} catch (error: any) {
// Precondition failed - concurrent modification
if (error.code === 412) {
if (attempt === maxRetries - 1) {
this.logger.error(`Max retries (${maxRetries}) exceeded for HNSW system data`)
throw new Error('Failed to save HNSW system data: max retries exceeded due to concurrent modifications')
}
const backoffMs = 50 * Math.pow(2, attempt)
await new Promise(resolve => setTimeout(resolve, backoffMs))
continue
}
// Other error - rethrow
this.logger.error('Failed to save HNSW system data:', error)
throw new Error(`Failed to save HNSW system data: ${error}`)
}
}
}
/**
* Get HNSW system data (entry point, max level)
* Storage path: system/hnsw-system.json
*/
public async getHNSWSystem(): Promise<{
entryPointId: string | null
maxLevel: number
} | null> {
await this.ensureInitialized()
try {
const key = `${this.systemPrefix}hnsw-system.json`
const file = this.bucket!.file(key)
const [contents] = await file.download()
return JSON.parse(contents.toString())
} catch (error: any) {
// GCS may return 404 in different formats depending on SDK version
const is404 =
error.code === 404 ||
error.statusCode === 404 ||
error.status === 404 ||
error.message?.includes('No such object') ||
error.message?.includes('not found') ||
error.message?.includes('404')
if (is404) {
return null
}
this.logger.error('Failed to get HNSW system data:', error)
throw new Error(`Failed to get HNSW system data: ${error}`)
}
}
// ============================================================================
// GCS Lifecycle Management & Autoclass
// Cost optimization through automatic tier transitions and Autoclass
// ============================================================================
/**
* Set lifecycle policy for automatic tier transitions and deletions
*
* GCS Storage Classes:
* - STANDARD: Hot data, most expensive (~$0.020/GB/month)
* - NEARLINE: <1 access/month (~$0.010/GB/month, 50% cheaper)
* - COLDLINE: <1 access/quarter (~$0.004/GB/month, 80% cheaper)
* - ARCHIVE: <1 access/year (~$0.0012/GB/month, 94% cheaper!)
*
* Example usage:
* ```typescript
* await storage.setLifecyclePolicy({
* rules: [
* {
* action: { type: 'SetStorageClass', storageClass: 'NEARLINE' },
* condition: { age: 30 }
* },
* {
* action: { type: 'SetStorageClass', storageClass: 'COLDLINE' },
* condition: { age: 90 }
* },
* {
* action: { type: 'Delete' },
* condition: { age: 365 }
* }
* ]
* })
* ```
*
* @param options Lifecycle configuration with rules for transitions and deletions
*/
public async setLifecyclePolicy(options: {
rules: Array<{
action: {
type: 'Delete' | 'SetStorageClass'
storageClass?: 'STANDARD' | 'NEARLINE' | 'COLDLINE' | 'ARCHIVE'
}
condition: {
age?: number // Days since object creation
createdBefore?: string // ISO 8601 date
matchesPrefix?: string[]
matchesSuffix?: string[]
}
}>
}): Promise<void> {
await this.ensureInitialized()
try {
this.logger.info(`Setting GCS lifecycle policy with ${options.rules.length} rules`)
// GCS lifecycle rules format
const lifecycleRules = options.rules.map(rule => {
const gcsRule: any = {
action: {
type: rule.action.type
},
condition: {}
}
// Add storage class for SetStorageClass action
if (rule.action.type === 'SetStorageClass' && rule.action.storageClass) {
gcsRule.action.storageClass = rule.action.storageClass
}
// Add conditions
if (rule.condition.age !== undefined) {
gcsRule.condition.age = rule.condition.age
}
if (rule.condition.createdBefore) {
gcsRule.condition.createdBefore = rule.condition.createdBefore
}
if (rule.condition.matchesPrefix) {
gcsRule.condition.matchesPrefix = rule.condition.matchesPrefix
}
if (rule.condition.matchesSuffix) {
gcsRule.condition.matchesSuffix = rule.condition.matchesSuffix
}
return gcsRule
})
// Update bucket lifecycle configuration
await this.bucket!.setMetadata({
lifecycle: {
rule: lifecycleRules
}
})
this.logger.info(`Successfully set lifecycle policy with ${options.rules.length} rules`)
} catch (error: any) {
this.logger.error('Failed to set lifecycle policy:', error)
throw new Error(`Failed to set GCS lifecycle policy: ${error.message || error}`)
}
}
/**
* Get current lifecycle policy configuration
*
* @returns Lifecycle configuration with all rules, or null if no policy is set
*/
public async getLifecyclePolicy(): Promise<{
rules: Array<{
action: {
type: string
storageClass?: string
}
condition: {
age?: number
createdBefore?: string
matchesPrefix?: string[]
matchesSuffix?: string[]
}
}>
} | null> {
await this.ensureInitialized()
try {
this.logger.info('Getting GCS lifecycle policy')
const [metadata] = await this.bucket!.getMetadata()
if (!metadata.lifecycle || !metadata.lifecycle.rule || metadata.lifecycle.rule.length === 0) {
this.logger.info('No lifecycle policy configured')
return null
}
// Convert GCS format to our format
const rules = metadata.lifecycle.rule.map((rule: any) => ({
action: {
type: rule.action.type,
...(rule.action.storageClass && { storageClass: rule.action.storageClass })
},
condition: {
...(rule.condition.age !== undefined && { age: rule.condition.age }),
...(rule.condition.createdBefore && { createdBefore: rule.condition.createdBefore }),
...(rule.condition.matchesPrefix && { matchesPrefix: rule.condition.matchesPrefix }),
...(rule.condition.matchesSuffix && { matchesSuffix: rule.condition.matchesSuffix })
}
}))
this.logger.info(`Found lifecycle policy with ${rules.length} rules`)
return { rules }
} catch (error: any) {
this.logger.error('Failed to get lifecycle policy:', error)
throw new Error(`Failed to get GCS lifecycle policy: ${error.message || error}`)
}
}
/**
* Remove lifecycle policy from bucket
*/
public async removeLifecyclePolicy(): Promise<void> {
await this.ensureInitialized()
try {
this.logger.info('Removing GCS lifecycle policy')
// Remove lifecycle configuration
await this.bucket!.setMetadata({
lifecycle: null
})
this.logger.info('Successfully removed lifecycle policy')
} catch (error: any) {
this.logger.error('Failed to remove lifecycle policy:', error)
throw new Error(`Failed to remove GCS lifecycle policy: ${error.message || error}`)
}
}
/**
* Enable Autoclass for automatic storage class optimization
*
* GCS Autoclass automatically moves objects between storage classes based on access patterns:
* - Frequent Access → STANDARD
* - Infrequent Access (30 days) → NEARLINE
* - Rarely Accessed (90 days) → COLDLINE
* - Archive Access (365 days) → ARCHIVE
*
* Benefits:
* - Automatic optimization based on access patterns (no manual rules needed)
* - No early deletion fees
* - No retrieval fees for NEARLINE/COLDLINE (only ARCHIVE has retrieval fees)
* - Up to 94% cost savings automatically
*
* Note: Autoclass is a bucket-level feature that requires bucket.update permission.
* It cannot be enabled per-object or per-prefix.
*
* @param options Autoclass configuration
*/
public async enableAutoclass(options: {
terminalStorageClass?: 'NEARLINE' | 'ARCHIVE' // Coldest storage class to use
} = {}): Promise<void> {
await this.ensureInitialized()
try {
this.logger.info('Enabling GCS Autoclass')
const autoclassConfig: any = {
enabled: true
}
// Set terminal storage class if specified
if (options.terminalStorageClass) {
autoclassConfig.terminalStorageClass = options.terminalStorageClass
}
await this.bucket!.setMetadata({
autoclass: autoclassConfig
})
this.logger.info(`Successfully enabled Autoclass${options.terminalStorageClass ? ` with terminal class ${options.terminalStorageClass}` : ''}`)
} catch (error: any) {
this.logger.error('Failed to enable Autoclass:', error)
throw new Error(`Failed to enable GCS Autoclass: ${error.message || error}`)
}
}
/**
* Get Autoclass configuration and status
*
* @returns Autoclass status, or null if not configured
*/
public async getAutoclassStatus(): Promise<{
enabled: boolean
terminalStorageClass?: string
toggleTime?: string
} | null> {
await this.ensureInitialized()
try {
this.logger.info('Getting GCS Autoclass status')
const [metadata] = await this.bucket!.getMetadata()
if (!metadata.autoclass) {
this.logger.info('Autoclass not configured')
return null
}
const status = {
enabled: metadata.autoclass.enabled || false,
...(metadata.autoclass.terminalStorageClass && {
terminalStorageClass: metadata.autoclass.terminalStorageClass
}),
...(metadata.autoclass.toggleTime && {
toggleTime: metadata.autoclass.toggleTime
})
}
this.logger.info(`Autoclass status: ${status.enabled ? 'enabled' : 'disabled'}`)
return status
} catch (error: any) {
this.logger.error('Failed to get Autoclass status:', error)
throw new Error(`Failed to get GCS Autoclass status: ${error.message || error}`)
}
}
/**
* Disable Autoclass for the bucket
*/
public async disableAutoclass(): Promise<void> {
await this.ensureInitialized()
try {
this.logger.info('Disabling GCS Autoclass')
await this.bucket!.setMetadata({
autoclass: {
enabled: false
}
})
this.logger.info('Successfully disabled Autoclass')
} catch (error: any) {
this.logger.error('Failed to disable Autoclass:', error)
throw new Error(`Failed to disable GCS Autoclass: ${error.message || error}`)
}
}
}