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}`, ); };