From adbe19aa16991c4dfd60213567d0e0f567c7e4ca Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Mon, 13 Apr 2026 15:36:07 -0400 Subject: [PATCH] fix: concurrent uploads, overwrite-in-place, 15min timeout - Bump Lambda timeout from 5min to 15min (9k+ POs need more time) - Upload S3 files 25x concurrently instead of sequentially - Replace clear-then-write with overwrite-in-place + delete stale to avoid S3 404s during concurrent KB ingestion jobs --- lambda/po-sync/index.ts | 90 ++++++++++++++++++++++++--------------- lib/constructs/po-sync.ts | 2 +- 2 files changed, 57 insertions(+), 35 deletions(-) diff --git a/lambda/po-sync/index.ts b/lambda/po-sync/index.ts index bc5fa21..be90186 100644 --- a/lambda/po-sync/index.ts +++ b/lambda/po-sync/index.ts @@ -133,9 +133,10 @@ async function scanAllPOs(): Promise[]> { return records; } -async function clearOldFiles(bucket: string): Promise { +/** List all existing S3 keys under the purchase-orders/ prefix. */ +async function listExistingKeys(bucket: string): Promise> { + const keys = new Set(); let continuationToken: string | undefined; - const toDelete: { Key: string }[] = []; do { const res = await s3.send( @@ -146,22 +147,27 @@ async function clearOldFiles(bucket: string): Promise { }), ); for (const obj of res.Contents ?? []) { - if (obj.Key) toDelete.push({ Key: obj.Key }); + if (obj.Key) keys.add(obj.Key); } continuationToken = res.NextContinuationToken; } while (continuationToken); - if (toDelete.length === 0) return; + return keys; +} - for (let i = 0; i < toDelete.length; i += 1000) { +/** Delete S3 keys that no longer correspond to active records. */ +async function deleteStaleKeys(bucket: string, keys: string[]): Promise { + if (keys.length === 0) return; + + for (let i = 0; i < keys.length; i += 1000) { await s3.send( new DeleteObjectsCommand({ Bucket: bucket, - Delete: { Objects: toDelete.slice(i, i + 1000) }, + Delete: { Objects: keys.slice(i, i + 1000).map((Key) => ({ Key })) }, }), ); } - console.log(`Deleted ${toDelete.length} stale PO file(s) from S3`); + console.log(`Deleted ${keys.length} stale PO file(s) from S3`); } export const handler = async (): Promise => { @@ -169,49 +175,65 @@ export const handler = async (): Promise => { const kbId = process.env.KNOWLEDGE_BASE_ID!; const dsId = process.env.DATA_SOURCE_ID!; - console.log('Clearing stale PO files from S3...'); - await clearOldFiles(bucket); + console.log('Listing existing PO files in S3...'); + const existingKeys = await listExistingKeys(bucket); + console.log(`Found ${existingKeys.size} existing file(s) in S3`); console.log('Scanning purchase-orders table...'); const pos = await scanAllPOs(); console.log(`Found ${pos.length} purchase order(s)`); + const writtenKeys = new Set(); let synced = 0; let failed = 0; - for (const po of pos) { - const poNumber = po.po_number; - if (!poNumber) { - console.warn('Skipping record with no po_number'); - failed++; - continue; - } + // Upload in batches of 25 concurrent requests + const CONCURRENCY = 25; + for (let i = 0; i < pos.length; i += CONCURRENCY) { + const batch = pos.slice(i, i + CONCURRENCY); + const results = await Promise.allSettled( + batch.map(async (po) => { + const poNumber = po.po_number; + if (!poNumber) { + console.warn('Skipping record with no po_number'); + throw new Error('no po_number'); + } - try { - const md = poToMarkdown(po); - // Sanitize PO number for use as S3 key - const safeKey = poNumber.replace(/[^a-zA-Z0-9._-]/g, '_'); + const md = poToMarkdown(po); + const safeKey = poNumber.replace(/[^a-zA-Z0-9._-]/g, '_'); + const key = `${S3_PREFIX}${safeKey}.md`; - await s3.send( - new PutObjectCommand({ - Bucket: bucket, - Key: `${S3_PREFIX}${safeKey}.md`, - Body: md, - ContentType: 'text/markdown', - Metadata: { 'po-number': poNumber }, - }), - ); + await s3.send( + new PutObjectCommand({ + Bucket: bucket, + Key: key, + Body: md, + ContentType: 'text/markdown', + Metadata: { 'po-number': poNumber }, + }), + ); - console.log(`Synced: PO ${poNumber}`); - synced++; - } catch (err) { - console.error(`Failed to sync PO ${poNumber}:`, err); - failed++; + return key; + }), + ); + + for (const r of results) { + if (r.status === 'fulfilled') { + writtenKeys.add(r.value); + synced++; + } else { + console.error(`Failed to sync PO:`, r.reason); + failed++; + } } } console.log(`Upload complete. ${synced} synced, ${failed} failed.`); + // Delete S3 files for records that no longer exist in DynamoDB + const staleKeys = [...existingKeys].filter((k) => !writtenKeys.has(k)); + await deleteStaleKeys(bucket, staleKeys); + console.log('Triggering Bedrock KB ingestion job...'); const ingestionRes = await bedrockAgent.send( new StartIngestionJobCommand({ diff --git a/lib/constructs/po-sync.ts b/lib/constructs/po-sync.ts index a44d8a0..5e6772a 100644 --- a/lib/constructs/po-sync.ts +++ b/lib/constructs/po-sync.ts @@ -33,7 +33,7 @@ export class PoSyncConstruct extends Construct { entry: path.join(__dirname, '../../lambda/po-sync/index.ts'), runtime: lambda.Runtime.NODEJS_22_X, memorySize: 512, - timeout: cdk.Duration.minutes(5), + timeout: cdk.Duration.minutes(15), environment: { PO_TABLE: poTable.tableName, KB_BUCKET_NAME: props.kbDocsBucket.bucketName,