diff --git a/lambda/po-sync/index.ts b/lambda/po-sync/index.ts new file mode 100644 index 0000000..be90186 --- /dev/null +++ b/lambda/po-sync/index.ts @@ -0,0 +1,247 @@ +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; +} + +/** List all existing S3 keys under the purchase-orders/ prefix. */ +async function listExistingKeys(bucket: string): Promise> { + const keys = new Set(); + let continuationToken: string | undefined; + + do { + const res = await s3.send( + new ListObjectsV2Command({ + Bucket: bucket, + Prefix: S3_PREFIX, + ContinuationToken: continuationToken, + }), + ); + for (const obj of res.Contents ?? []) { + if (obj.Key) keys.add(obj.Key); + } + continuationToken = res.NextContinuationToken; + } while (continuationToken); + + return keys; +} + +/** 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: keys.slice(i, i + 1000).map((Key) => ({ Key })) }, + }), + ); + } + console.log(`Deleted ${keys.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('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; + + // 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'); + } + + 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: key, + Body: md, + ContentType: 'text/markdown', + Metadata: { 'po-number': poNumber }, + }), + ); + + 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({ + 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..5e6772a --- /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(15), + 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,