Merge pull request #7 from Sea-Haven-Industries/feature/workorder-kb-sync

feat: work orders → Bedrock KB sync
This commit is contained in:
Adam Moussa 2026-04-13 15:38:20 -04:00 • committed by GitHub
commit 1dc275dee4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 353 additions and 0 deletions

View file

@ -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<Set<string>> {
const keys = new Set<string>();
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<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 work order file(s) from S3`);
}
async function scanAllWorkOrders(): Promise<WorkOrder[]> {
const results: WorkOrder[] = [];
let lastKey: Record<string, any> | 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<WorkOrderComment[]> {
const results: WorkOrderComment[] = [];
let lastKey: Record<string, any> | 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<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 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<string>();
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}`,
);
};

View file

@ -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));
}
}

View file

@ -6,6 +6,7 @@ import { BedrockAgentConstruct } from './constructs/bedrock-agent';
import { SlackHandlerConstruct } from './constructs/slack-handler'; import { SlackHandlerConstruct } from './constructs/slack-handler';
import { NotionSyncConstruct } from './constructs/notion-sync'; import { NotionSyncConstruct } from './constructs/notion-sync';
import { PoSyncConstruct } from './constructs/po-sync'; import { PoSyncConstruct } from './constructs/po-sync';
import { WorkorderSyncConstruct } from './constructs/workorder-sync';
export class SeahavenSlackBotStack extends cdk.Stack { export class SeahavenSlackBotStack extends cdk.Stack {
constructor(scope: Construct, id: string, props?: cdk.StackProps) { constructor(scope: Construct, id: string, props?: cdk.StackProps) {
@ -53,6 +54,14 @@ export class SeahavenSlackBotStack extends cdk.Stack {
dataSourceId: knowledgeBase.dataSource.dataSourceId, 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) ─────────────────────────── // ── Slack webhook handler (API Gateway + Lambda) ───────────────────────────
new SlackHandlerConstruct(this, 'SlackHandler', { new SlackHandlerConstruct(this, 'SlackHandler', {
accountId: this.account, accountId: this.account,