From 615970eba296a22dea2d63ae17e1757aebe67ed0 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Mon, 13 Apr 2026 14:33:53 -0400 Subject: [PATCH 1/2] =?UTF-8?q?feat:=20daily=20purchase-orders=20DynamoDB?= =?UTF-8?q?=20=E2=86=92=20Bedrock=20KB=20sync?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add a new Lambda and CDK construct that scans the purchase-orders DynamoDB table (owned by po-ingest), converts each PO to markdown, uploads to S3 under the purchase-orders/ prefix, and triggers a Bedrock Knowledge Base ingestion job. Runs daily at 02:00 UTC via EventBridge alongside the existing Notion sync. --- lambda/po-sync/index.ts | 225 ++++++++++++++++++++++++++++++++ lib/constructs/po-sync.ts | 71 ++++++++++ lib/seahaven-slack-bot-stack.ts | 9 ++ 3 files changed, 305 insertions(+) create mode 100644 lambda/po-sync/index.ts create mode 100644 lib/constructs/po-sync.ts diff --git a/lambda/po-sync/index.ts b/lambda/po-sync/index.ts new file mode 100644 index 0000000..bc5fa21 --- /dev/null +++ b/lambda/po-sync/index.ts @@ -0,0 +1,225 @@ +import { + DynamoDBClient, + ScanCommand, + type AttributeValue, +} from '@aws-sdk/client-dynamodb'; +import { + S3Client, + PutObjectCommand, + ListObjectsV2Command, + DeleteObjectsCommand, +} from '@aws-sdk/client-s3'; +import { + BedrockAgentClient, + StartIngestionJobCommand, +} from '@aws-sdk/client-bedrock-agent'; + +const dynamo = new DynamoDBClient({}); +const s3 = new S3Client({}); +const bedrockAgent = new BedrockAgentClient({ region: process.env.REGION! }); + +const S3_PREFIX = 'purchase-orders/'; + +/** Unwrap a DynamoDB AttributeValue to a plain JS value. */ +function unwrap(attr: AttributeValue | undefined): any { + if (!attr) return undefined; + if (attr.S !== undefined) return attr.S; + if (attr.N !== undefined) return Number(attr.N); + if (attr.BOOL !== undefined) return attr.BOOL; + if (attr.NULL) return null; + if (attr.L) return attr.L.map(unwrap); + if (attr.M) { + const obj: Record = {}; + for (const [k, v] of Object.entries(attr.M)) { + obj[k] = unwrap(v); + } + return obj; + } + return undefined; +} + +/** Convert a raw DynamoDB item to a plain JS object. */ +function itemToRecord(item: Record): Record { + const rec: Record = {}; + for (const [k, v] of Object.entries(item)) { + rec[k] = unwrap(v); + } + return rec; +} + +/** Render a PO record as a human-readable markdown document. */ +function poToMarkdown(po: Record): string { + const lines: string[] = []; + + lines.push(`# Purchase Order ${po.po_number}`); + lines.push(''); + + if (po.po_status) lines.push(`**Status:** ${po.po_status}`); + if (po.email_type === 'cancellation') lines.push(`**Cancelled:** ${po.cancelled_at ?? 'yes'}`); + if (po.source_system) lines.push(`**Source:** ${po.source_system}`); + if (po.order_date) lines.push(`**Order Date:** ${po.order_date}`); + if (po.revision_date) lines.push(`**Revision Date:** ${po.revision_date}`); + if (po.payment_terms) lines.push(`**Payment Terms:** ${po.payment_terms}`); + if (po.requisition_number) lines.push(`**Requisition #:** ${po.requisition_number}`); + if (po.department) lines.push(`**Department:** ${po.department}`); + if (po.submitted_by) lines.push(`**Submitted By:** ${po.submitted_by}`); + if (po.on_behalf_of) lines.push(`**On Behalf Of:** ${po.on_behalf_of}`); + + if (po.total_amount !== undefined) { + lines.push(`**Total:** ${po.currency ?? 'USD'} ${po.total_amount}`); + } + + // Supplier + if (po.supplier?.name) { + lines.push(''); + lines.push(`## Supplier`); + lines.push(`**Name:** ${po.supplier.name}`); + } + + // Ship-to + const ship = po.ship_to; + if (ship) { + lines.push(''); + lines.push(`## Ship To`); + if (ship.name) lines.push(`**Name:** ${ship.name}`); + if (ship.address) lines.push(`**Address:** ${ship.address}`); + if (ship.location_code) lines.push(`**Location Code:** ${ship.location_code}`); + if (ship.attn) lines.push(`**Attn:** ${ship.attn}`); + } + + // Line items + const items = po.line_items; + if (Array.isArray(items) && items.length > 0) { + lines.push(''); + lines.push('## Line Items'); + lines.push(''); + lines.push('| Description | Amount | Currency | Need By | Category | Account Code | Period |'); + lines.push('|---|---|---|---|---|---|---|'); + for (const li of items) { + const row = [ + li.description ?? '', + li.amount ?? '', + li.currency ?? '', + li.need_by ?? '', + li.category ?? '', + li.account_code ?? '', + li.period ?? '', + ]; + lines.push(`| ${row.join(' | ')} |`); + } + } + + lines.push(''); + return lines.join('\n'); +} + +async function scanAllPOs(): Promise[]> { + const records: Record[] = []; + let lastKey: Record | undefined; + + do { + const res = await dynamo.send( + new ScanCommand({ + TableName: process.env.PO_TABLE!, + ExclusiveStartKey: lastKey, + }), + ); + for (const item of res.Items ?? []) { + records.push(itemToRecord(item)); + } + lastKey = res.LastEvaluatedKey; + } while (lastKey); + + return records; +} + +async function clearOldFiles(bucket: string): Promise { + let continuationToken: string | undefined; + const toDelete: { Key: string }[] = []; + + do { + const res = await s3.send( + new ListObjectsV2Command({ + Bucket: bucket, + Prefix: S3_PREFIX, + ContinuationToken: continuationToken, + }), + ); + for (const obj of res.Contents ?? []) { + if (obj.Key) toDelete.push({ Key: obj.Key }); + } + continuationToken = res.NextContinuationToken; + } while (continuationToken); + + if (toDelete.length === 0) return; + + for (let i = 0; i < toDelete.length; i += 1000) { + await s3.send( + new DeleteObjectsCommand({ + Bucket: bucket, + Delete: { Objects: toDelete.slice(i, i + 1000) }, + }), + ); + } + console.log(`Deleted ${toDelete.length} stale PO file(s) from S3`); +} + +export const handler = async (): Promise => { + const bucket = process.env.KB_BUCKET_NAME!; + 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('Scanning purchase-orders table...'); + const pos = await scanAllPOs(); + console.log(`Found ${pos.length} purchase order(s)`); + + 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; + } + + try { + const md = poToMarkdown(po); + // Sanitize PO number for use as S3 key + const safeKey = poNumber.replace(/[^a-zA-Z0-9._-]/g, '_'); + + await s3.send( + new PutObjectCommand({ + Bucket: bucket, + Key: `${S3_PREFIX}${safeKey}.md`, + 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++; + } + } + + console.log(`Upload complete. ${synced} synced, ${failed} failed.`); + + console.log('Triggering Bedrock KB ingestion job...'); + const ingestionRes = await bedrockAgent.send( + new StartIngestionJobCommand({ + knowledgeBaseId: kbId, + dataSourceId: dsId, + }), + ); + console.log( + `Ingestion job started: ${ingestionRes.ingestionJob?.ingestionJobId}`, + ); +}; diff --git a/lib/constructs/po-sync.ts b/lib/constructs/po-sync.ts new file mode 100644 index 0000000..a44d8a0 --- /dev/null +++ b/lib/constructs/po-sync.ts @@ -0,0 +1,71 @@ +import * as cdk from 'aws-cdk-lib'; +import { Construct } from 'constructs'; +import * as lambda from 'aws-cdk-lib/aws-lambda'; +import * as lambdaNodejs from 'aws-cdk-lib/aws-lambda-nodejs'; +import * as s3 from 'aws-cdk-lib/aws-s3'; +import * as dynamodb from 'aws-cdk-lib/aws-dynamodb'; +import * as events from 'aws-cdk-lib/aws-events'; +import * as targets from 'aws-cdk-lib/aws-events-targets'; +import * as iam from 'aws-cdk-lib/aws-iam'; +import * as path from 'path'; + +export interface PoSyncProps { + region: string; + kbDocsBucket: s3.Bucket; + knowledgeBaseId: string; + dataSourceId: string; +} + +export class PoSyncConstruct extends Construct { + public readonly syncLambda: lambdaNodejs.NodejsFunction; + + constructor(scope: Construct, id: string, props: PoSyncProps) { + super(scope, id); + + // Import the purchase-orders DynamoDB table (owned by po-ingest stack) + const poTable = dynamodb.Table.fromTableName( + this, 'PurchaseOrdersTable', 'purchase-orders', + ); + + // Lambda — scans purchase-orders table, uploads markdown to S3, triggers KB ingestion + this.syncLambda = new lambdaNodejs.NodejsFunction(this, 'SyncLambda', { + functionName: 'seahaven-po-sync', + entry: path.join(__dirname, '../../lambda/po-sync/index.ts'), + runtime: lambda.Runtime.NODEJS_22_X, + memorySize: 512, + timeout: cdk.Duration.minutes(5), + environment: { + PO_TABLE: poTable.tableName, + KB_BUCKET_NAME: props.kbDocsBucket.bucketName, + KNOWLEDGE_BASE_ID: props.knowledgeBaseId, + DATA_SOURCE_ID: props.dataSourceId, + REGION: props.region, + }, + }); + + // Grant read access to the purchase-orders table + poTable.grantReadData(this.syncLambda); + + // Grant read/write to the KB docs bucket (read to list+delete old, write new) + props.kbDocsBucket.grantReadWrite(this.syncLambda); + + // Allow Lambda to start a Bedrock KB ingestion job + this.syncLambda.addToRolePolicy( + new iam.PolicyStatement({ + actions: ['bedrock:StartIngestionJob'], + resources: [ + `arn:aws:bedrock:${props.region}:*:knowledge-base/${props.knowledgeBaseId}`, + ], + }), + ); + + // EventBridge rule — fires daily at 02:00 UTC + const dailyRule = new events.Rule(this, 'DailySyncRule', { + ruleName: 'seahaven-po-daily-sync', + description: 'Daily purchase-orders → KB sync at 02:00 UTC', + schedule: events.Schedule.cron({ minute: '0', hour: '2' }), + }); + + dailyRule.addTarget(new targets.LambdaFunction(this.syncLambda)); + } +} diff --git a/lib/seahaven-slack-bot-stack.ts b/lib/seahaven-slack-bot-stack.ts index f5cbca0..daa0d82 100644 --- a/lib/seahaven-slack-bot-stack.ts +++ b/lib/seahaven-slack-bot-stack.ts @@ -5,6 +5,7 @@ import { KnowledgeBaseConstruct } from './constructs/knowledge-base'; import { BedrockAgentConstruct } from './constructs/bedrock-agent'; import { SlackHandlerConstruct } from './constructs/slack-handler'; import { NotionSyncConstruct } from './constructs/notion-sync'; +import { PoSyncConstruct } from './constructs/po-sync'; export class SeahavenSlackBotStack extends cdk.Stack { constructor(scope: Construct, id: string, props?: cdk.StackProps) { @@ -44,6 +45,14 @@ export class SeahavenSlackBotStack extends cdk.Stack { dataSourceId: knowledgeBase.dataSource.dataSourceId, }); + // ── Purchase Orders → KB daily sync (EventBridge + Lambda) ──────────────── + new PoSyncConstruct(this, 'PoSync', { + region: this.region, + kbDocsBucket: knowledgeBase.docsBucket, + knowledgeBaseId: knowledgeBase.knowledgeBase.knowledgeBaseId, + dataSourceId: knowledgeBase.dataSource.dataSourceId, + }); + // ── Slack webhook handler (API Gateway + Lambda) ─────────────────────────── new SlackHandlerConstruct(this, 'SlackHandler', { accountId: this.account, -- 2.50.1 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 2/2] 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, -- 2.50.1