diff --git a/package-lock.json b/package-lock.json index bb1d81e5..7bb9e815 100644 --- a/package-lock.json +++ b/package-lock.json @@ -11,6 +11,7 @@ "dependencies": { "@aws-sdk/client-s3": "^3.540.0", "@huggingface/transformers": "^3.1.0", + "@smithy/node-http-handler": "^4.1.1", "buffer": "^6.0.3", "dotenv": "^16.4.5", "uuid": "^9.0.1" @@ -2926,12 +2927,12 @@ ] }, "node_modules/@smithy/abort-controller": { - "version": "4.0.4", - "resolved": "https://registry.npmjs.org/@smithy/abort-controller/-/abort-controller-4.0.4.tgz", - "integrity": "sha512-gJnEjZMvigPDQWHrW3oPrFhQtkrgqBkyjj3pCIdF3A5M6vsZODG93KNlfJprv6bp4245bdT32fsHK4kkH3KYDA==", + "version": "4.0.5", + "resolved": "https://registry.npmjs.org/@smithy/abort-controller/-/abort-controller-4.0.5.tgz", + "integrity": "sha512-jcrqdTQurIrBbUm4W2YdLVMQDoL0sA9DTxYd2s+R/y+2U9NLOP7Xf/YqfSg1FZhlZIYEnvk2mwbyvIfdLEPo8g==", "license": "Apache-2.0", "dependencies": { - "@smithy/types": "^4.3.1", + "@smithy/types": "^4.3.2", "tslib": "^2.6.2" }, "engines": { @@ -3280,15 +3281,15 @@ } }, "node_modules/@smithy/node-http-handler": { - "version": "4.1.0", - "resolved": "https://registry.npmjs.org/@smithy/node-http-handler/-/node-http-handler-4.1.0.tgz", - "integrity": "sha512-vqfSiHz2v8b3TTTrdXi03vNz1KLYYS3bhHCDv36FYDqxT7jvTll1mMnCrkD+gOvgwybuunh/2VmvOMqwBegxEg==", + "version": "4.1.1", + "resolved": "https://registry.npmjs.org/@smithy/node-http-handler/-/node-http-handler-4.1.1.tgz", + "integrity": "sha512-RHnlHqFpoVdjSPPiYy/t40Zovf3BBHc2oemgD7VsVTFFZrU5erFFe0n52OANZZ/5sbshgD93sOh5r6I35Xmpaw==", "license": "Apache-2.0", "dependencies": { - "@smithy/abort-controller": "^4.0.4", - "@smithy/protocol-http": "^5.1.2", - "@smithy/querystring-builder": "^4.0.4", - "@smithy/types": "^4.3.1", + "@smithy/abort-controller": "^4.0.5", + "@smithy/protocol-http": "^5.1.3", + "@smithy/querystring-builder": "^4.0.5", + "@smithy/types": "^4.3.2", "tslib": "^2.6.2" }, "engines": { @@ -3309,12 +3310,12 @@ } }, "node_modules/@smithy/protocol-http": { - "version": "5.1.2", - "resolved": "https://registry.npmjs.org/@smithy/protocol-http/-/protocol-http-5.1.2.tgz", - "integrity": "sha512-rOG5cNLBXovxIrICSBm95dLqzfvxjEmuZx4KK3hWwPFHGdW3lxY0fZNXfv2zebfRO7sJZ5pKJYHScsqopeIWtQ==", + "version": "5.1.3", + "resolved": "https://registry.npmjs.org/@smithy/protocol-http/-/protocol-http-5.1.3.tgz", + "integrity": "sha512-fCJd2ZR7D22XhDY0l+92pUag/7je2BztPRQ01gU5bMChcyI0rlly7QFibnYHzcxDvccMjlpM/Q1ev8ceRIb48w==", "license": "Apache-2.0", "dependencies": { - "@smithy/types": "^4.3.1", + "@smithy/types": "^4.3.2", "tslib": "^2.6.2" }, "engines": { @@ -3322,12 +3323,12 @@ } }, "node_modules/@smithy/querystring-builder": { - "version": "4.0.4", - "resolved": "https://registry.npmjs.org/@smithy/querystring-builder/-/querystring-builder-4.0.4.tgz", - "integrity": "sha512-SwREZcDnEYoh9tLNgMbpop+UTGq44Hl9tdj3rf+yeLcfH7+J8OXEBaMc2kDxtyRHu8BhSg9ADEx0gFHvpJgU8w==", + "version": "4.0.5", + "resolved": "https://registry.npmjs.org/@smithy/querystring-builder/-/querystring-builder-4.0.5.tgz", + "integrity": "sha512-NJeSCU57piZ56c+/wY+AbAw6rxCCAOZLCIniRE7wqvndqxcKKDOXzwWjrY7wGKEISfhL9gBbAaWWgHsUGedk+A==", "license": "Apache-2.0", "dependencies": { - "@smithy/types": "^4.3.1", + "@smithy/types": "^4.3.2", "@smithy/util-uri-escape": "^4.0.0", "tslib": "^2.6.2" }, @@ -3411,9 +3412,9 @@ } }, "node_modules/@smithy/types": { - "version": "4.3.1", - "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.3.1.tgz", - "integrity": "sha512-UqKOQBL2x6+HWl3P+3QqFD4ncKq0I8Nuz9QItGv5WuKuMHuuwlhvqcZCoXGfc+P1QmfJE7VieykoYYmrOoFJxA==", + "version": "4.3.2", + "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.3.2.tgz", + "integrity": "sha512-QO4zghLxiQ5W9UZmX2Lo0nta2PuE1sSrXUYDoaB6HMR762C0P7v/HEPHf6ZdglTVssJG1bsrSBxdc3quvDSihw==", "license": "Apache-2.0", "dependencies": { "tslib": "^2.6.2" diff --git a/package.json b/package.json index 3816dee1..b1842368 100644 --- a/package.json +++ b/package.json @@ -156,6 +156,7 @@ "dependencies": { "@aws-sdk/client-s3": "^3.540.0", "@huggingface/transformers": "^3.1.0", + "@smithy/node-http-handler": "^4.1.1", "buffer": "^6.0.3", "dotenv": "^16.4.5", "uuid": "^9.0.1" diff --git a/src/storage/adapters/s3CompatibleStorage.ts b/src/storage/adapters/s3CompatibleStorage.ts index 1040b330..d79e1a69 100644 --- a/src/storage/adapters/s3CompatibleStorage.ts +++ b/src/storage/adapters/s3CompatibleStorage.ts @@ -97,6 +97,16 @@ export class S3CompatibleStorage extends BaseStorage { // Change log for efficient synchronization private changeLogPrefix: string = 'change-log/' + // 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() + // Operation executors for timeout and retry handling private operationExecutors: StorageOperationExecutors @@ -168,6 +178,17 @@ export class S3CompatibleStorage extends BaseStorage { try { // Import AWS SDK modules only when needed const { S3Client } = await import('@aws-sdk/client-s3') + const { NodeHttpHandler } = await import('@smithy/node-http-handler') + const https = await import('https') + + // Configure HTTP agent with high socket limits for high-volume scenarios + // This prevents socket exhaustion without requiring user configuration + const httpAgent = new https.Agent({ + keepAlive: true, + maxSockets: 500, // Increased from default 50 to handle high volume + maxFreeSockets: 100, // Keep more sockets alive for reuse + timeout: 60000 // 60 second timeout + }) // Configure the S3 client based on the service type const clientConfig: any = { @@ -175,7 +196,16 @@ export class S3CompatibleStorage extends BaseStorage { credentials: { accessKeyId: this.accessKeyId, secretAccessKey: this.secretAccessKey - } + }, + // Use custom HTTP handler with optimized socket management + requestHandler: new NodeHttpHandler({ + httpsAgent: httpAgent, + connectionTimeout: 10000, // 10 second connection timeout + socketTimeout: 60000 // 60 second socket timeout + }), + // Retry configuration for resilience + maxAttempts: 5, // Retry up to 5 times + retryMode: 'adaptive' // Use adaptive retry with backoff } // Add session token if provided @@ -212,7 +242,7 @@ export class S3CompatibleStorage extends BaseStorage { getMany: async (ids: string[]) => { const result = new Map() // Process in batches to avoid overwhelming the S3 API - const batchSize = 10 + const batchSize = this.getBatchSize() const batches: string[][] = [] // Split into batches @@ -253,7 +283,7 @@ export class S3CompatibleStorage extends BaseStorage { getMany: async (ids: string[]) => { const result = new Map() // Process in batches to avoid overwhelming the S3 API - const batchSize = 10 + const batchSize = this.getBatchSize() const batches: string[][] = [] // Split into batches @@ -301,6 +331,91 @@ export class S3CompatibleStorage extends BaseStorage { } } + /** + * Dynamically adjust batch size based on memory pressure and error rates + */ + private adjustBatchSize(): void { + const now = Date.now() + + // Check memory pressure periodically + if (now - this.lastMemoryCheck > this.memoryCheckInterval) { + this.lastMemoryCheck = now + const memUsage = process.memoryUsage() + const heapUsedPercent = (memUsage.heapUsed / memUsage.heapTotal) * 100 + + // Adjust batch size based on memory pressure + if (heapUsedPercent > 80) { + // High memory pressure - reduce batch size + this.currentBatchSize = Math.max(1, Math.floor(this.currentBatchSize * 0.5)) + this.maxConcurrentOperations = Math.max(10, Math.floor(this.maxConcurrentOperations * 0.5)) + this.logger.warn(`High memory pressure (${heapUsedPercent.toFixed(1)}%), reducing batch size to ${this.currentBatchSize}`) + } else if (heapUsedPercent < 50 && this.consecutiveErrors === 0) { + // Low memory pressure and no errors - gradually increase batch size + this.currentBatchSize = Math.min(this.baseBatchSize * 2, this.currentBatchSize + 1) + this.maxConcurrentOperations = Math.min(200, this.maxConcurrentOperations + 10) + } + } + + // Reset error counter periodically if no recent errors + if (now - this.lastErrorReset > 60000 && this.consecutiveErrors > 0) { + this.consecutiveErrors = Math.max(0, this.consecutiveErrors - 1) + this.lastErrorReset = now + } + + // Adjust based on error rate + if (this.consecutiveErrors > 5) { + this.currentBatchSize = Math.max(1, Math.floor(this.currentBatchSize * 0.7)) + this.maxConcurrentOperations = Math.max(10, Math.floor(this.maxConcurrentOperations * 0.7)) + this.logger.warn(`High error rate, reducing batch size to ${this.currentBatchSize}`) + } + } + + /** + * Apply backpressure when system is under load + */ + private async applyBackpressure(): Promise { + // Wait if too many operations are pending + while (this.pendingOperations >= this.maxConcurrentOperations) { + const memUsage = process.memoryUsage() + const heapUsedPercent = (memUsage.heapUsed / memUsage.heapTotal) * 100 + + if (heapUsedPercent > 85) { + // Critical memory pressure - wait longer + await new Promise(resolve => setTimeout(resolve, 100)) + } else { + // Normal backpressure - short wait + await new Promise(resolve => setImmediate(resolve)) + } + } + + this.pendingOperations++ + } + + /** + * Release backpressure after operation completes + */ + private releaseBackpressure(success: boolean = true): void { + this.pendingOperations = Math.max(0, this.pendingOperations - 1) + + if (!success) { + this.consecutiveErrors++ + } else if (this.consecutiveErrors > 0) { + // Gradually reduce error count on success + this.consecutiveErrors = Math.max(0, this.consecutiveErrors - 0.5) + } + + // Adjust batch size based on current conditions + this.adjustBatchSize() + } + + /** + * Get current batch size for operations + */ + private getBatchSize(): number { + this.adjustBatchSize() + return this.currentBatchSize + } + /** * Save a noun to storage (internal implementation) */ @@ -314,6 +429,9 @@ export class S3CompatibleStorage extends BaseStorage { protected async saveNode(node: HNSWNode): Promise { await this.ensureInitialized() + // Apply backpressure before starting operation + await this.applyBackpressure() + try { this.logger.trace(`Saving node ${node.id}`) @@ -380,7 +498,11 @@ export class S3CompatibleStorage extends BaseStorage { verifyError ) } + // Release backpressure on success + this.releaseBackpressure(true) } catch (error) { + // Release backpressure on error + this.releaseBackpressure(false) this.logger.error(`Failed to save node ${node.id}:`, error) throw new Error(`Failed to save node ${node.id}: ${error}`) } @@ -726,6 +848,9 @@ export class S3CompatibleStorage extends BaseStorage { protected async saveEdge(edge: Edge): Promise { await this.ensureInitialized() + // Apply backpressure before starting operation + await this.applyBackpressure() + try { // Convert connections Map to a serializable format const serializableEdge = { @@ -758,7 +883,12 @@ export class S3CompatibleStorage extends BaseStorage { vector: edge.vector } }) + + // Release backpressure on success + this.releaseBackpressure(true) } catch (error) { + // Release backpressure on error + this.releaseBackpressure(false) this.logger.error(`Failed to save edge ${edge.id}:`, error) throw new Error(`Failed to save edge ${edge.id}: ${error}`) } @@ -1197,6 +1327,9 @@ export class S3CompatibleStorage extends BaseStorage { public async saveMetadata(id: string, metadata: any): Promise { await this.ensureInitialized() + // Apply backpressure before starting operation + await this.applyBackpressure() + try { // Import the PutObjectCommand only when needed const { PutObjectCommand } = await import('@aws-sdk/client-s3') @@ -1250,7 +1383,12 @@ export class S3CompatibleStorage extends BaseStorage { verifyError ) } + + // Release backpressure on success + this.releaseBackpressure(true) } catch (error) { + // Release backpressure on error + this.releaseBackpressure(false) this.logger.error(`Failed to save metadata for ${id}:`, error) throw new Error(`Failed to save metadata for ${id}: ${error}`) } diff --git a/tests/mocks/s3-mock.ts b/tests/mocks/s3-mock.ts index 1078b305..0a28b347 100644 --- a/tests/mocks/s3-mock.ts +++ b/tests/mocks/s3-mock.ts @@ -175,8 +175,8 @@ function handlePutObject(command: any) { const parsedBody = JSON.parse(bodyContent) console.log(`Parsed JSON body:`, parsedBody) - // If this is a noun or verb, ensure it has an id property - if (Key.includes('/nouns/') || Key.includes('/verbs/')) { + // If this is a noun or verb (but NOT metadata), ensure it has an id property + if ((Key.includes('/nouns/') || Key.includes('/verbs/')) && !Key.includes('/metadata/')) { if (!parsedBody.id) { console.error(`Warning: Object ${Key} does not have an id property`) // Add id property based on the key name @@ -255,8 +255,8 @@ function handleGetObject(command: any) { const parsedBody = JSON.parse(bodyContent) console.log(`Parsed JSON body for ${Key}:`, parsedBody) - // If this is a noun or verb, ensure it has an id property - if (Key.includes('/nouns/') || Key.includes('/verbs/')) { + // If this is a noun or verb (but NOT metadata), ensure it has an id property + if ((Key.includes('/nouns/') || Key.includes('/verbs/')) && !Key.includes('/metadata/')) { if (!parsedBody.id) { console.error(`Warning: Object ${Key} does not have an id property`) // Add id property based on the key name