From 72401477361e9be145563cb7480c6b3f79efff1f Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Mon, 13 Apr 2026 14:35:47 -0400 Subject: [PATCH 1/2] =?UTF-8?q?feat:=20daily=20work=20orders=20DynamoDB=20?= =?UTF-8?q?=E2=86=92=20Bedrock=20KB=20sync?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add a new Lambda and CDK construct that scans the WorkOrders and WorkOrderComments DynamoDB tables (owned by workorder-ingest), converts each work order + comment history to markdown, uploads to S3 under the work-orders/ prefix, and triggers a Bedrock Knowledge Base ingestion job. Runs daily at 02:00 UTC via EventBridge alongside the existing Notion sync. --- lambda/workorder-sync/index.ts | 245 +++++++++++++++++++++++++++++++ lib/constructs/workorder-sync.ts | 76 ++++++++++ lib/seahaven-slack-bot-stack.ts | 9 ++ 3 files changed, 330 insertions(+) create mode 100644 lambda/workorder-sync/index.ts create mode 100644 lib/constructs/workorder-sync.ts diff --git a/lambda/workorder-sync/index.ts b/lambda/workorder-sync/index.ts new file mode 100644 index 0000000..b18340a --- /dev/null +++ b/lambda/workorder-sync/index.ts @@ -0,0 +1,245 @@ +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; +} + +async function clearOldFiles(bucket: string): Promise { + let continuationToken: string | undefined; + const toDelete: { Key: string }[] = []; + + do { + const res = await s3.send( + new ListObjectsV2Command({ + Bucket: bucket, + Prefix: WO_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; + + // S3 DeleteObjects accepts up to 1000 keys per request + 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 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('Clearing stale work order files from S3...'); + await clearOldFiles(bucket); + + console.log('Scanning work orders from DynamoDB...'); + const workOrders = await scanAllWorkOrders(); + console.log(`Found ${workOrders.length} work order(s)`); + + let synced = 0; + let failed = 0; + + for (const wo of workOrders) { + try { + const comments = await getComments(wo.work_order_id); + const markdown = formatWorkOrderMarkdown(wo, comments); + + await s3.send( + new PutObjectCommand({ + Bucket: bucket, + Key: `${WO_PREFIX}${wo.work_order_id}.md`, + Body: markdown, + ContentType: 'text/markdown', + Metadata: { + 'work-order-id': wo.work_order_id, + 'site-code': (wo.site_code ?? '').slice(0, 256), + }, + }), + ); + + console.log(`Synced: ${wo.work_order_id} (${comments.length} comment(s))`); + synced++; + } catch (err) { + console.error(`Failed to sync work order ${wo.work_order_id}:`, 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}`, + ); +}; 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 f5cbca0..302c288 100644 --- a/lib/seahaven-slack-bot-stack.ts +++ b/lib/seahaven-slack-bot-stack.ts @@ -5,6 +5,7 @@ import { KnowledgeBaseConstruct } from './constructs/knowledge-base'; import { BedrockAgentConstruct } from './constructs/bedrock-agent'; import { SlackHandlerConstruct } from './constructs/slack-handler'; import { NotionSyncConstruct } from './constructs/notion-sync'; +import { WorkorderSyncConstruct } from './constructs/workorder-sync'; export class SeahavenSlackBotStack extends cdk.Stack { constructor(scope: Construct, id: string, props?: cdk.StackProps) { @@ -44,6 +45,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, From 641494530b36b8dc7138c6ae0f321b6cd5397df8 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Mon, 13 Apr 2026 15:36:16 -0400 Subject: [PATCH 2/2] fix: concurrent uploads and overwrite-in-place sync strategy - Upload S3 files 10x concurrently instead of sequentially - Replace clear-then-write with overwrite-in-place + delete stale to avoid S3 404s during concurrent KB ingestion jobs --- lambda/workorder-sync/index.ts | 85 +++++++++++++++++++++------------- 1 file changed, 54 insertions(+), 31 deletions(-) diff --git a/lambda/workorder-sync/index.ts b/lambda/workorder-sync/index.ts index b18340a..45100ab 100644 --- a/lambda/workorder-sync/index.ts +++ b/lambda/workorder-sync/index.ts @@ -52,9 +52,10 @@ interface WorkOrderComment { ingested_at?: string; } -async function clearOldFiles(bucket: string): Promise { +/** List all existing S3 keys under the work-orders/ prefix. */ +async function listExistingKeys(bucket: string): Promise> { + const keys = new Set(); let continuationToken: string | undefined; - const toDelete: { Key: string }[] = []; do { const res = await s3.send( @@ -65,23 +66,27 @@ async function clearOldFiles(bucket: string): Promise { }), ); for (const obj of res.Contents ?? []) { - if (obj.Key) toDelete.push({ Key: obj.Key }); + if (obj.Key) keys.add(obj.Key); } continuationToken = res.NextContinuationToken; } while (continuationToken); - if (toDelete.length === 0) return; + return keys; +} - // S3 DeleteObjects accepts up to 1000 keys per request - for (let i = 0; i < toDelete.length; i += 1000) { +/** 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: toDelete.slice(i, i + 1000) }, + Delete: { Objects: keys.slice(i, i + 1000).map((Key) => ({ Key })) }, }), ); } - console.log(`Deleted ${toDelete.length} stale work order file(s) from S3`); + console.log(`Deleted ${keys.length} stale work order file(s) from S3`); } async function scanAllWorkOrders(): Promise { @@ -194,44 +199,62 @@ export const handler = async (): Promise => { const kbId = process.env.KNOWLEDGE_BASE_ID!; const dsId = process.env.DATA_SOURCE_ID!; - console.log('Clearing stale work order files from S3...'); - await clearOldFiles(bucket); + 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; - for (const wo of workOrders) { - try { - const comments = await getComments(wo.work_order_id); - const markdown = formatWorkOrderMarkdown(wo, comments); + // 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: `${WO_PREFIX}${wo.work_order_id}.md`, - Body: markdown, - ContentType: 'text/markdown', - Metadata: { - 'work-order-id': wo.work_order_id, - 'site-code': (wo.site_code ?? '').slice(0, 256), - }, - }), - ); + 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), + }, + }), + ); - console.log(`Synced: ${wo.work_order_id} (${comments.length} comment(s))`); - synced++; - } catch (err) { - console.error(`Failed to sync work order ${wo.work_order_id}:`, err); - failed++; + 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({