- 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
247 lines
7.3 KiB
TypeScript
247 lines
7.3 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;
|
|
}
|
|
|
|
/** 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}`,
|
|
);
|
|
};
|