import { S3Client, PutObjectCommand, ListObjectsV2Command, DeleteObjectsCommand, } from '@aws-sdk/client-s3'; import { DynamoDBClient, ScanCommand, QueryCommand, } from '@aws-sdk/client-dynamodb'; import { unmarshall } from '@aws-sdk/util-dynamodb'; import { BedrockAgentClient, StartIngestionJobCommand, } from '@aws-sdk/client-bedrock-agent'; const s3 = new S3Client({}); const dynamo = new DynamoDBClient({}); const bedrockAgent = new BedrockAgentClient({ region: process.env.REGION! }); const WO_PREFIX = 'work-orders/'; interface WorkOrder { work_order_id: string; description?: string; wo_status?: string; customer?: string; site_code?: string; building?: string; address?: string; severity?: string; priority?: string; date_reported?: string; scheduled_start?: string; due_date?: string; assigned_to?: string; source_email_s3_key?: string; created_at?: string; updated_at?: string; record_type?: string; } interface WorkOrderComment { work_order_id: string; comment_id: string; record_type?: string; commenter?: string; text?: string; created_at?: string; source_email_s3_key?: string; ingested_at?: string; } /** List all existing S3 keys under the work-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: WO_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 work order file(s) from S3`); } async function scanAllWorkOrders(): Promise { const results: WorkOrder[] = []; let lastKey: Record | undefined; do { const res = await dynamo.send( new ScanCommand({ TableName: process.env.WORK_ORDERS_TABLE!, ExclusiveStartKey: lastKey, }), ); for (const item of res.Items ?? []) { results.push(unmarshall(item) as WorkOrder); } lastKey = res.LastEvaluatedKey; } while (lastKey); return results; } async function getComments(workOrderId: string): Promise { const results: WorkOrderComment[] = []; let lastKey: Record | undefined; do { const res = await dynamo.send( new QueryCommand({ TableName: process.env.COMMENTS_TABLE!, KeyConditionExpression: 'work_order_id = :woid', ExpressionAttributeValues: { ':woid': { S: workOrderId }, }, ExclusiveStartKey: lastKey, }), ); for (const item of res.Items ?? []) { results.push(unmarshall(item) as WorkOrderComment); } lastKey = res.LastEvaluatedKey; } while (lastKey); // Sort chronologically by created_at results.sort((a, b) => (a.created_at ?? '').localeCompare(b.created_at ?? '')); return results; } function formatWorkOrderMarkdown(wo: WorkOrder, comments: WorkOrderComment[]): string { const lines: string[] = []; lines.push(`# Work Order: ${wo.work_order_id}`); lines.push(''); if (wo.description) { lines.push(`## Description`); lines.push(''); lines.push(wo.description); lines.push(''); } lines.push(`## Details`); lines.push(''); lines.push(`| Field | Value |`); lines.push(`| ----- | ----- |`); if (wo.wo_status) lines.push(`| Status | ${wo.wo_status} |`); if (wo.customer) lines.push(`| Customer | ${wo.customer} |`); if (wo.site_code) lines.push(`| Site Code | ${wo.site_code} |`); if (wo.building) lines.push(`| Building | ${wo.building} |`); if (wo.address) lines.push(`| Address | ${wo.address} |`); if (wo.severity) lines.push(`| Severity | ${wo.severity} |`); if (wo.priority) lines.push(`| Priority | ${wo.priority} |`); if (wo.assigned_to) lines.push(`| Assigned To | ${wo.assigned_to} |`); if (wo.record_type) lines.push(`| Record Type | ${wo.record_type} |`); lines.push(''); lines.push(`## Dates`); lines.push(''); lines.push(`| Field | Value |`); lines.push(`| ----- | ----- |`); if (wo.date_reported) lines.push(`| Date Reported | ${wo.date_reported} |`); if (wo.scheduled_start) lines.push(`| Scheduled Start | ${wo.scheduled_start} |`); if (wo.due_date) lines.push(`| Due Date | ${wo.due_date} |`); if (wo.created_at) lines.push(`| Created At | ${wo.created_at} |`); if (wo.updated_at) lines.push(`| Updated At | ${wo.updated_at} |`); lines.push(''); if (comments.length > 0) { lines.push(`## Comment History`); lines.push(''); for (const comment of comments) { const timestamp = comment.created_at ?? 'unknown date'; const author = comment.commenter ?? 'unknown'; const type = comment.record_type ? ` [${comment.record_type}]` : ''; lines.push(`### ${timestamp} - ${author}${type}`); lines.push(''); if (comment.text) { lines.push(comment.text); lines.push(''); } } } return lines.join('\n'); } 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 work order files in S3...'); const existingKeys = await listExistingKeys(bucket); console.log(`Found ${existingKeys.size} existing file(s) in S3`); console.log('Scanning work orders from DynamoDB...'); const workOrders = await scanAllWorkOrders(); console.log(`Found ${workOrders.length} work order(s)`); const writtenKeys = new Set(); let synced = 0; let failed = 0; // Upload in batches of 10 (each WO also queries comments, so keep concurrency moderate) const CONCURRENCY = 10; for (let i = 0; i < workOrders.length; i += CONCURRENCY) { const batch = workOrders.slice(i, i + CONCURRENCY); const results = await Promise.allSettled( batch.map(async (wo) => { const comments = await getComments(wo.work_order_id); const markdown = formatWorkOrderMarkdown(wo, comments); const key = `${WO_PREFIX}${wo.work_order_id}.md`; await s3.send( new PutObjectCommand({ Bucket: bucket, Key: key, Body: markdown, ContentType: 'text/markdown', Metadata: { 'work-order-id': wo.work_order_id, 'site-code': (wo.site_code ?? '').slice(0, 256), }, }), ); return key; }), ); for (const r of results) { if (r.status === 'fulfilled') { writtenKeys.add(r.value); synced++; } else { console.error(`Failed to sync work order:`, 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}`, ); };