This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
seahaven-slack-bot/lambda/po-sync/index.ts
Adam Moussa 615970eba2 feat: daily purchase-orders DynamoDB → Bedrock KB sync
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.
2026-04-13 14:33:53 -04:00

225 lines
6.4 KiB
TypeScript

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;
}
async function clearOldFiles(bucket: string): Promise<void> {
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<void> => {
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}`,
);
};