diff --git a/lambda/workorder-sync/index.ts b/lambda/workorder-sync/index.ts new file mode 100644 index 0000000..45100ab --- /dev/null +++ b/lambda/workorder-sync/index.ts @@ -0,0 +1,268 @@ +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}`, + ); +}; diff --git a/lib/constructs/workorder-sync.ts b/lib/constructs/workorder-sync.ts new file mode 100644 index 0000000..eb28fd5 --- /dev/null +++ b/lib/constructs/workorder-sync.ts @@ -0,0 +1,76 @@ +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 WorkorderSyncProps { + region: string; + kbDocsBucket: s3.Bucket; + knowledgeBaseId: string; + dataSourceId: string; +} + +export class WorkorderSyncConstruct extends Construct { + public readonly syncLambda: lambdaNodejs.NodejsFunction; + + constructor(scope: Construct, id: string, props: WorkorderSyncProps) { + super(scope, id); + + // Import existing DynamoDB tables by name (owned by workorder-ingest stack) + const workOrdersTable = dynamodb.Table.fromTableName( + this, 'WorkOrdersTable', 'WorkOrders', + ); + const commentsTable = dynamodb.Table.fromTableName( + this, 'WorkOrderCommentsTable', 'WorkOrderComments', + ); + + // Lambda — scans DynamoDB work orders + comments, uploads markdown to S3, triggers KB ingestion + this.syncLambda = new lambdaNodejs.NodejsFunction(this, 'SyncLambda', { + functionName: 'seahaven-workorder-sync', + entry: path.join(__dirname, '../../lambda/workorder-sync/index.ts'), + runtime: lambda.Runtime.NODEJS_22_X, + memorySize: 512, + timeout: cdk.Duration.minutes(5), + environment: { + WORK_ORDERS_TABLE: workOrdersTable.tableName, + COMMENTS_TABLE: commentsTable.tableName, + KB_BUCKET_NAME: props.kbDocsBucket.bucketName, + KNOWLEDGE_BASE_ID: props.knowledgeBaseId, + DATA_SOURCE_ID: props.dataSourceId, + REGION: props.region, + }, + }); + + // Grant read access to both DynamoDB tables + workOrdersTable.grantReadData(this.syncLambda); + commentsTable.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-workorder-daily-sync', + description: 'Daily Work Orders → KB sync at 02:00 UTC', + schedule: events.Schedule.cron({ minute: '0', hour: '2' }), + }); + + dailyRule.addTarget(new targets.LambdaFunction(this.syncLambda)); + } +} diff --git a/lib/seahaven-slack-bot-stack.ts b/lib/seahaven-slack-bot-stack.ts index daa0d82..d7814f2 100644 --- a/lib/seahaven-slack-bot-stack.ts +++ b/lib/seahaven-slack-bot-stack.ts @@ -6,6 +6,7 @@ import { BedrockAgentConstruct } from './constructs/bedrock-agent'; import { SlackHandlerConstruct } from './constructs/slack-handler'; import { NotionSyncConstruct } from './constructs/notion-sync'; import { PoSyncConstruct } from './constructs/po-sync'; +import { WorkorderSyncConstruct } from './constructs/workorder-sync'; export class SeahavenSlackBotStack extends cdk.Stack { constructor(scope: Construct, id: string, props?: cdk.StackProps) { @@ -53,6 +54,14 @@ export class SeahavenSlackBotStack extends cdk.Stack { dataSourceId: knowledgeBase.dataSource.dataSourceId, }); + // ── Work Orders → KB daily sync (EventBridge + Lambda) ──────────────────── + new WorkorderSyncConstruct(this, 'WorkorderSync', { + 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,