From d1d04dc92e58cba04dfe35f3900bce8373c6f524 Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 16:57:37 -0300 Subject: [PATCH 1/6] fix(AAIDomains): keyset pagination + per-page splitting to fix OOM on large syncs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Replace offset .range() pagination with cursor-based keyset (.gt('id', lastId)) eliminates O(n^2) Supabase scans on large tables - Move textSplitter.splitDocuments() inside load() per page instead of splitting the entire 10k doc array after load() — eliminates dual large arrays in heap simultaneously - Add textSplitter param to AAIDomainsLoaderParams and AAIDomainsLoader - Fixes KUMELLO-68: OOM crashes when syncing 10k+ domains on Render Standard (512MB) Other loaders unaffected — change is isolated to AAIDomains.ts only. --- .../documentloaders/AAIDomains/AAIDomains.ts | 44 +++++++++++-------- 1 file changed, 25 insertions(+), 19 deletions(-) diff --git a/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts b/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts index 3df934eb6cc..c3ed8fee90b 100644 --- a/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts +++ b/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts @@ -261,16 +261,9 @@ class AAIDomains_DocumentLoaders implements INode { contentFields: parsedContentFields } - const loader = new AAIDomainsLoader(loaderOptions) + const loader = new AAIDomainsLoader({ ...loaderOptions, textSplitter: textSplitter || null }) - let docs: IDocument[] = [] - - if (textSplitter) { - docs = await loader.load() - docs = await textSplitter.splitDocuments(docs) - } else { - docs = await loader.load() - } + const docs = await loader.load() // Apply metadata const parsedMetadata = metadata ? (typeof metadata === 'object' ? metadata : JSON.parse(metadata)) : null @@ -305,6 +298,7 @@ interface AAIDomainsLoaderParams { isValid: string | null hasAnalysis: string contentFields: string[] | null + textSplitter?: TextSplitter | null } class AAIDomainsLoader extends BaseDocumentLoader { @@ -318,6 +312,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { private isValid: string | null private hasAnalysis: string private contentFields: string[] | null + private textSplitter: TextSplitter | null constructor(params: AAIDomainsLoaderParams) { super() @@ -331,6 +326,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { this.isValid = params.isValid this.hasAnalysis = params.hasAnalysis this.contentFields = params.contentFields + this.textSplitter = params.textSplitter ?? null } public async load(): Promise { @@ -350,10 +346,10 @@ class AAIDomainsLoader extends BaseDocumentLoader { } }) - // Use larger page size since we're only selecting essential fields const pageSize = Math.min(this.limit, 100) let allDocs: IDocument[] = [] - let currentPage = 0 + let lastId: string | null = null + let pageNum = 0 console.info('[AAIDomains] Starting load with params:', { limit: this.limit, @@ -369,7 +365,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { const remainingItems = this.limit - allDocs.length const currentPageSize = Math.min(pageSize, remainingItems) - console.info(`[AAIDomains] Fetching page ${currentPage}, size ${currentPageSize}`) + console.info(`[AAIDomains] Fetching page ${pageNum}, size ${currentPageSize}`) // Retry logic with exponential backoff let retryCount = 0 @@ -446,8 +442,12 @@ class AAIDomainsLoader extends BaseDocumentLoader { let query = supabase .from('domains') .select(selectFields) - .order('updated_at', { ascending: false }) - .range(currentPage * pageSize, (currentPage + 1) * pageSize - 1) + .order('id', { ascending: true }) + .limit(currentPageSize) + + if (lastId) { + query = query.gt('id', lastId) + } // Apply filters if (this.searchTerm) { @@ -479,7 +479,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { throw new Error(`Failed to fetch domains from AAI Datastore: ${error.message}`) } - console.info(`[AAIDomains] Page ${currentPage} response:`, { + console.info(`[AAIDomains] Page ${pageNum} response:`, { hasData: !!data, domainsCount: data?.length || 0 }) @@ -490,6 +490,9 @@ class AAIDomainsLoader extends BaseDocumentLoader { break } + lastId = (data as any[])[data.length - 1].id + pageNum++ + // Transform domain_tags array to flat tags array (non-mutating) const domainsWithTags = (data as any[]).map((domain: any) => ({ ...domain, @@ -503,16 +506,19 @@ class AAIDomainsLoader extends BaseDocumentLoader { } const pageDocs = filteredDomains.map((d: any) => this.createDocumentFromDomain(d)) - allDocs.push(...pageDocs) - currentPage++ - // Stop if we've fetched enough + if (this.textSplitter) { + const pageChunks = await this.textSplitter.splitDocuments(pageDocs) + allDocs.push(...pageChunks) + } else { + allDocs.push(...pageDocs) + } + if (allDocs.length >= this.limit || data.length < currentPageSize) { shouldStopPagination = true break } - // Success - break out of retry loop break } catch (error: any) { lastError = error From eac7c8099d6e4fa7f874f50a9a1a7286dfdfd236 Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 17:33:35 -0300 Subject: [PATCH 2/6] fix(documentstore): replace Promise.all individual saves with bulk insert in syncAndRefreshChunks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 500 concurrent chunkRepository.save() calls per batch were opening 500 simultaneous DB connections to Supabase, causing connection pool exhaustion and service crashes on large document stores (10k+ chunks). Replace with single chunkRepository.insert(entities) per batch — same pattern already used in _saveChunksToStorage. One SQL INSERT per batch instead of 500 individual transactions. Also adds sanitizeChunkContent() call to strip null bytes, consistent with the _saveChunksToStorage implementation. --- .../src/services/documentstore/index.ts | 57 +++++++------------ 1 file changed, 22 insertions(+), 35 deletions(-) diff --git a/packages/server/src/services/documentstore/index.ts b/packages/server/src/services/documentstore/index.ts index e69861988a8..54764dac009 100644 --- a/packages/server/src/services/documentstore/index.ts +++ b/packages/server/src/services/documentstore/index.ts @@ -464,23 +464,17 @@ const syncAndRefreshChunks = async (storeId: string, fileId: string, userId: str for (let i = 0; i < docs.length; i += SAVE_BATCH_SIZE) { const batch = docs.slice(i, i + SAVE_BATCH_SIZE) try { - await Promise.all( - batch.map(async (chunk: IDocument, localIndex: number) => { - const globalIndex = i + localIndex - const docChunk: DocumentStoreFileChunk = { - userId, - organizationId, - docId: fileId, - storeId: storeId, - id: uuidv4(), - chunkNo: globalIndex + 1, - pageContent: chunk.pageContent, - metadata: JSON.stringify(chunk.metadata) - } - const dChunk = chunkRepository.create(docChunk) - await chunkRepository.save(dChunk) - }) - ) + const entities = batch.map((chunk: IDocument, localIndex: number) => ({ + userId, + organizationId, + docId: fileId, + storeId: storeId, + id: uuidv4(), + chunkNo: i + localIndex + 1, + pageContent: sanitizeChunkContent(chunk.pageContent), + metadata: JSON.stringify(chunk.metadata) + })) + await chunkRepository.insert(entities) persistedChunks += batch.length persistedChars += batch.reduce((acc: number, chunk: IDocument) => acc + (chunk.pageContent?.length ?? 0), 0) // Free memory: allow GC to reclaim saved chunks @@ -1300,26 +1294,19 @@ const _saveChunksToStorage = async ( for (let i = 0; i < docs.length; i += SAVE_BATCH_SIZE) { const batch = docs.slice(i, i + SAVE_BATCH_SIZE) try { - await Promise.all( - batch.map(async (chunk: IDocument, localIndex: number) => { - const globalIndex = i + localIndex - const docChunk: DocumentStoreFileChunk = { - docId: newLoaderId, - storeId: data.storeId || '', - id: uuidv4(), - chunkNo: globalIndex + 1, - pageContent: sanitizeChunkContent(chunk.pageContent), - metadata: JSON.stringify(chunk.metadata), - userId: data.userId, - organizationId: data.organizationId - } - const dChunk = chunkRepository.create(docChunk) - await chunkRepository.save(dChunk) - }) - ) + const entities = batch.map((chunk: IDocument, localIndex: number) => ({ + docId: newLoaderId, + storeId: data.storeId || '', + id: uuidv4(), + chunkNo: i + localIndex + 1, + pageContent: sanitizeChunkContent(chunk.pageContent), + metadata: JSON.stringify(chunk.metadata), + userId: data.userId, + organizationId: data.organizationId + })) + await chunkRepository.insert(entities) persistedChunks += batch.length persistedChars += batch.reduce((acc: number, chunk: IDocument) => acc + (chunk.pageContent?.length ?? 0), 0) - // Free memory: allow GC to reclaim saved chunks for (let j = i; j < i + batch.length; j++) { docs[j] = null as any } From dc9f18619451a835a8f4cf454f1c82916b817e78 Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 17:34:19 -0300 Subject: [PATCH 3/6] fix(docker): set Node.js heap to 384MB to match Render Standard 512MB container NODE_OPTIONS was set to --max-old-space-size=8192 (8GB) but the Render Standard plan only provides 512MB RAM. V8 GC wouldn't collect aggressively because it believed it had 8GB of headroom, causing the OS to OOM-kill the process at 512MB with no stack trace. Setting to 384MB (512MB - ~128MB OS/overhead) forces V8 to GC proactively and throw a JS heap error instead of a silent OOM kill. --- Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Dockerfile b/Dockerfile index 25ade7ea2ca..98173eb91d3 100644 --- a/Dockerfile +++ b/Dockerfile @@ -29,7 +29,7 @@ RUN pnpm config set store-dir /root/.pnpm-store ENV PUPPETEER_SKIP_DOWNLOAD=true ENV PUPPETEER_EXECUTABLE_PATH=/usr/bin/chromium-browser -ENV NODE_OPTIONS=--max-old-space-size=8192 +ENV NODE_OPTIONS=--max-old-space-size=384 ################################################################################ # Prune projects From eaaefdd0351f0f4a5b4ba33198147dad77f0dd3b Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 17:40:46 -0300 Subject: [PATCH 4/6] fix(docker): raise NODE_OPTIONS heap to 4096MB to stop OOM crashes --- Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Dockerfile b/Dockerfile index 98173eb91d3..5d3be512a63 100644 --- a/Dockerfile +++ b/Dockerfile @@ -29,7 +29,7 @@ RUN pnpm config set store-dir /root/.pnpm-store ENV PUPPETEER_SKIP_DOWNLOAD=true ENV PUPPETEER_EXECUTABLE_PATH=/usr/bin/chromium-browser -ENV NODE_OPTIONS=--max-old-space-size=384 +ENV NODE_OPTIONS=--max-old-space-size=4096 ################################################################################ # Prune projects From 1ce1227be0670f04381a175495c31d2ad55e15da Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 17:42:29 -0300 Subject: [PATCH 5/6] docs(documentstore): annotate SAVE_BATCH_SIZE with PostgreSQL parameter ceiling --- packages/server/src/services/documentstore/index.ts | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/server/src/services/documentstore/index.ts b/packages/server/src/services/documentstore/index.ts index 54764dac009..c212f938135 100644 --- a/packages/server/src/services/documentstore/index.ts +++ b/packages/server/src/services/documentstore/index.ts @@ -58,6 +58,9 @@ import { Telemetry } from '../../utils/telemetry' import nodesService from '../nodes' // Batch sizes for chunk DB operations and vector store upsert +// PostgreSQL parameter limit is 65535. DocumentStoreFileChunk has 8 columns, so +// SAVE_BATCH_SIZE=500 uses 4000 parameters (6% of limit) — safe headroom. +// If the entity gains columns or this value is raised, recalculate: rows × columns < 65535. const SAVE_BATCH_SIZE = 500 const UPSERT_BATCH_SIZE = 500 From 9ad510652b812cceb4761f5ebd2598da7c4e2a1b Mon Sep 17 00:00:00 2001 From: Max Techera Date: Mon, 16 Mar 2026 17:59:03 -0300 Subject: [PATCH 6/6] fix(AAIDomains): track source domain count separately from chunk count MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - while loop now guards on fetchedDomainCount (source docs) not allDocs.length (chunks), so limit=N always means N source domains regardless of whether a textSplitter is active - remainingItems calculation uses fetchedDomainCount, keeping page fetch size correctly calibrated to remaining source documents - lastId cursor advanced only after docs are successfully pushed to allDocs, preventing silent data gaps on transient errors mid-page - remove allDocs.slice(0, this.limit) — no longer correct when chunks > domains - note insert() lifecycle-hook bypass at both call sites in documentstore service --- .../documentloaders/AAIDomains/AAIDomains.ts | 23 +++++++++---------- .../src/services/documentstore/index.ts | 4 ++-- 2 files changed, 13 insertions(+), 14 deletions(-) diff --git a/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts b/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts index c3ed8fee90b..8f04e358c39 100644 --- a/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts +++ b/packages/components/nodes/documentloaders/AAIDomains/AAIDomains.ts @@ -348,6 +348,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { const pageSize = Math.min(this.limit, 100) let allDocs: IDocument[] = [] + let fetchedDomainCount = 0 let lastId: string | null = null let pageNum = 0 @@ -361,8 +362,8 @@ class AAIDomainsLoader extends BaseDocumentLoader { hasAnalysis: this.hasAnalysis }) - while (allDocs.length < this.limit) { - const remainingItems = this.limit - allDocs.length + while (fetchedDomainCount < this.limit) { + const remainingItems = this.limit - fetchedDomainCount const currentPageSize = Math.min(pageSize, remainingItems) console.info(`[AAIDomains] Fetching page ${pageNum}, size ${currentPageSize}`) @@ -490,9 +491,6 @@ class AAIDomainsLoader extends BaseDocumentLoader { break } - lastId = (data as any[])[data.length - 1].id - pageNum++ - // Transform domain_tags array to flat tags array (non-mutating) const domainsWithTags = (data as any[]).map((domain: any) => ({ ...domain, @@ -505,6 +503,8 @@ class AAIDomainsLoader extends BaseDocumentLoader { filteredDomains = this.filterByTags(domainsWithTags) } + fetchedDomainCount += filteredDomains.length + const pageDocs = filteredDomains.map((d: any) => this.createDocumentFromDomain(d)) if (this.textSplitter) { @@ -514,7 +514,11 @@ class AAIDomainsLoader extends BaseDocumentLoader { allDocs.push(...pageDocs) } - if (allDocs.length >= this.limit || data.length < currentPageSize) { + // Advance cursor only after docs are successfully pushed + lastId = (data as any[])[data.length - 1].id + pageNum++ + + if (fetchedDomainCount >= this.limit || data.length < currentPageSize) { shouldStopPagination = true break } @@ -547,12 +551,7 @@ class AAIDomainsLoader extends BaseDocumentLoader { await new Promise((resolve) => setTimeout(resolve, 200)) } - console.info(`[AAIDomains] Load complete. Total documents: ${allDocs.length}`) - - // Truncate to exact limit - if (allDocs.length > this.limit) { - allDocs = allDocs.slice(0, this.limit) - } + console.info(`[AAIDomains] Load complete. Total documents: ${allDocs.length}, source domains fetched: ${fetchedDomainCount}`) return allDocs } diff --git a/packages/server/src/services/documentstore/index.ts b/packages/server/src/services/documentstore/index.ts index c212f938135..75c0e19d62f 100644 --- a/packages/server/src/services/documentstore/index.ts +++ b/packages/server/src/services/documentstore/index.ts @@ -477,7 +477,7 @@ const syncAndRefreshChunks = async (storeId: string, fileId: string, userId: str pageContent: sanitizeChunkContent(chunk.pageContent), metadata: JSON.stringify(chunk.metadata) })) - await chunkRepository.insert(entities) + await chunkRepository.insert(entities) // insert() skips TypeORM lifecycle hooks — safe: DocumentStoreFileChunk has none persistedChunks += batch.length persistedChars += batch.reduce((acc: number, chunk: IDocument) => acc + (chunk.pageContent?.length ?? 0), 0) // Free memory: allow GC to reclaim saved chunks @@ -1307,7 +1307,7 @@ const _saveChunksToStorage = async ( userId: data.userId, organizationId: data.organizationId })) - await chunkRepository.insert(entities) + await chunkRepository.insert(entities) // insert() skips TypeORM lifecycle hooks — safe: DocumentStoreFileChunk has none persistedChunks += batch.length persistedChars += batch.reduce((acc: number, chunk: IDocument) => acc + (chunk.pageContent?.length ?? 0), 0) for (let j = i; j < i + batch.length; j++) {