feat: purchase orders → Bedrock KB sync #6

Merged
amoussa1229 merged 2 commits from feature/po-kb-sync into master 2026-04-13 19:37:29 +00:00
3 changed files with 327 additions and 0 deletions

247
lambda/po-sync/index.ts Normal file
View file

@ -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<string, any> = {};
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<string, AttributeValue>): Record<string, any> {
const rec: Record<string, any> = {};
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, any>): 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<Record<string, any>[]> {
const records: Record<string, any>[] = [];
let lastKey: Record<string, AttributeValue> | 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<Set<string>> {
const keys = new Set<string>();
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<void> {
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<void> => {
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<string>();
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}`,
);
};

71
lib/constructs/po-sync.ts Normal file
View file

@ -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));
}
}

View file

@ -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,