feat: Alex rollout — persona, Socket Mode, payments, channels, App Home, unanswered questions #9

Merged
amoussa1229 merged 3 commits from feature/alex-rollout into main 2026-04-30 22:57:23 +00:00
22 changed files with 2678 additions and 354 deletions

3
.gitignore vendored
View file

@ -22,3 +22,6 @@ docs/
# Claude Code project settings
.claude/
# macOS
.DS_Store

162
README.md
View file

@ -1,28 +1,31 @@
# Seahaven Slack Bot
# Seahaven Slack Bot — Alex
Internal Slack DM assistant for Sea Haven Industries, powered by AWS Bedrock. Employees can ask questions about vendors, company policies, SOPs, SA8000 compliance, work orders, purchase orders, and Amazon site assignments via direct message.
Internal Slack assistant for Sea Haven Industries, powered by AWS Bedrock. Employees can DM Alex or @mention Alex in any channel to ask about vendors, payments, work orders, purchase orders, Amazon sites, company policies, and more.
## Architecture
```
Slack DM
Slack (WebSocket)
│
▼
API Gateway (bot.seahaven.com)
ECS Fargate (Socket Mode) ← persistent WebSocket connection to Slack
│
├── POST /slack/events
│ ▼
│ slack-webhook Lambda ← verifies Slack signature, returns 200 immediately
│ │ (async invoke)
│ ▼
│ slack-processor Lambda ← calls Bedrock Agent, logs to DynamoDB, posts reply
│ │
│ ▼
│ Bedrock Agent (Claude Sonnet 4.5)
│ ├── Knowledge Base (AOSS + S3) ← SA8000 docs, SOPs, employee handbook, Notion pages, POs, work orders, site list
│ ├── QBO_Lookup action group ← QuickBooks vendor search (VPC, static IP)
│ ├── Google_Maps_Lookup action group ← fallback vendor search
│ └── WO_PO_Lookup action group ← work order, purchase order, and site lookups (DynamoDB direct)
├── DM (message.im) ──────────────────┐
├── @mention (app_mention) ────────────┤
│ async invoke
│ ▼
│ slack-processor Lambda ← calls Bedrock Agent, logs to DynamoDB, posts reply
│ │
│ ▼
│ Bedrock Agent — Alex (Claude Sonnet 4.5)
│ ├── Knowledge Base (AOSS + S3)
│ ├── QBO_Lookup action group ← QuickBooks vendor search (VPC, static IP)
│ ├── Google_Maps_Lookup action group ← fallback vendor search
│ └── WO_PO_Lookup action group ← WO, PO, site, and payment lookups (DynamoDB direct)
│
└── App Home opened ──────────────────→ app-home Lambda ← publishes Block Kit capabilities view
API Gateway (bot.seahaven.com)
│
└── GET /qbo/*
▼
@ -36,6 +39,15 @@ EventBridge (daily 02:00 UTC)
└── workorder-sync Lambda ← exports WorkOrders + Comments DynamoDB → S3 → KB ingestion
```
### What Alex Can Do
- **Vendor Lookup** — search QBO for existing vendors, KB for approved lists, Google Maps for new options
- **Payment & Invoice Status** — check if a vendor has been paid, if an invoice is scheduled, look up by check number
- **Work Order Lookups** — status, assignments, comments, and full history by WO number
- **Purchase Order Lookups** — status, line items, supplier, and ship-to details by PO number
- **Site Lookups** — Amazon facility address and details by site code, or list all sites in a state
- **Company Policies & SOPs** — employee handbook, SA8000 compliance, and operational procedures
### Vendor Query Priority
1. **QBO** — existing vetted vendors in QuickBooks Online
2. **Knowledge Base** — approved vendor lists in internal docs
@ -50,30 +62,33 @@ EventBridge (daily 02:00 UTC)
| LLM | Claude Sonnet 4.5 (`anthropic.claude-sonnet-4-5-20250929-v1:0`) |
| Embeddings | Amazon Titan Embed Text V2 (1024 dimensions) |
| Vector Store | OpenSearch Serverless (VECTORSEARCH) |
| Socket Mode | ECS Fargate (`seahaven-socket-mode`) — persistent WebSocket to Slack |
| Conversation Log | DynamoDB `seahaven-conversations` (90-day TTL) |
| Unanswered Questions | DynamoDB `seahaven-unanswered-questions` (180-day TTL) |
| Payment Data | DynamoDB `PaymentsDashboard` (via payments-dashboard) |
| PO Data Source | DynamoDB `purchase-orders` (via po-ingest) |
| Work Order Data Source | DynamoDB `WorkOrders` + `WorkOrderComments` (via workorder-ingest) |
| Site Assignments | DynamoDB `SiteAssignments` (seeded from CSV via `scripts/seed-sites.ts`) |
| VPC | `seahaven-vpc` (`vpc-0d3d4b67bd0cf8a68`) — QBO Lambdas only |
| Site Assignments | DynamoDB `verified-sites` (auto-populated via po-ingest Streams pipeline) |
| VPC | `seahaven-vpc` (`vpc-0d3d4b67bd0cf8a68`) — QBO Lambdas + Socket Mode |
| Static Outbound IP | `52.202.83.13` (NAT Gateway for Intuit IP allowlist) |
| Webhook URL | `https://bot.seahaven.com/slack/events` |
| QBO OAuth URLs | `/qbo/connect`, `/qbo/callback`, `/qbo/disconnect`, `/qbo/launch` |
| QBO OAuth URLs | `bot.seahaven.com/qbo/connect`, `/qbo/callback`, `/qbo/disconnect`, `/qbo/launch` |
## Prerequisites
- Node.js 22+
- AWS CDK (`npm install -g aws-cdk`)
- Docker Desktop (required by `@cdklabs/generative-ai-cdk-constructs` for Lambda bundling)
- Docker Desktop (required for Lambda bundling and Socket Mode container build)
- AWS credentials configured locally
## Secrets Manager
These secrets must exist before deploying. The Slack, QBO, and Maps secrets must be created manually before first deploy. The Notion secret is created automatically by CDK with a placeholder — update it after deploy.
These secrets must exist before deploying. The Slack, QBO, Maps, and app-level token secrets must be created manually before first deploy. The Notion secret is created automatically by CDK with a placeholder.
| Secret Name | Created | Structure |
|---|---|---|
| `seahaven/slack/credentials` | Manual (pre-deploy) | `{ "botToken": "xoxb-...", "signingSecret": "..." }` |
| `seahaven/qbo/oauth` | Manual (pre-deploy) | `{ "clientId": "", "clientSecret": "", "refreshToken": "", "realmId": "" }` — refreshToken and realmId are auto-populated via `/qbo/connect` OAuth flow |
| `seahaven/slack/app-level-token` | Manual (pre-deploy) | Plaintext `xapp-...` token with `connections:write` scope |
| `seahaven/qbo/oauth` | Manual (pre-deploy) | `{ "clientId": "", "clientSecret": "", "refreshToken": "", "realmId": "" }` |
| `seahaven/google/maps-api-key` | Manual (pre-deploy) | `{ "apiKey": "" }` |
| `seahaven/notion/api-key` | Auto (CDK) | `{ "apiKey": "secret_..." }` |
@ -85,19 +100,17 @@ npx cdk bootstrap aws://328440206208/us-east-1 # first time only
npx cdk deploy
```
The deploy takes ~15 minutes on first run. The AOSS collection and vector index creation are the slow steps — both handled automatically.
The deploy takes ~15 minutes on first run. The AOSS collection, vector index creation, and Docker image build are the slow steps.
Stack outputs after deploy:
- `KBDocsBucketName` — S3 bucket to upload knowledge base documents
- `AgentId` — Bedrock Agent ID
- `SlackWebhookUrl` — URL to register in Slack app settings
- `NotionSecretName` — Secrets Manager secret to populate with your Notion integration token
- `AgentId` — Bedrock Agent ID (Alex)
## Post-Deployment Setup
### 1. Upload knowledge base documents
Upload SA8000 compliance docs, SOPs, and the employee handbook to the S3 bucket printed in stack outputs. Supported formats: PDF, DOCX, TXT, HTML, CSV, XLSX.
Upload SA8000 compliance docs, SOPs, and the employee handbook to the S3 bucket printed in stack outputs.
```bash
aws s3 cp ./your-docs/ s3://<KBDocsBucketName>/ --recursive
@ -107,42 +120,30 @@ aws s3 cp ./your-docs/ s3://<KBDocsBucketName>/ --recursive
```bash
aws bedrock-agent start-ingestion-job \
--knowledge-base-id <AgentId from outputs> \
--knowledge-base-id <from Bedrock console> \
--data-source-id <DataSourceId> \
--region us-east-1
```
Or trigger from the Bedrock console: **Knowledge Bases → seahaven-kb → Sync**.
### 3. Configure Notion sync
1. Create a Notion internal integration at **Settings → Connections → Develop or manage integrations**
2. Name it `seahaven-office-ops` and associate it with the **Office Operations** teamspace
3. Copy the integration token (`secret_...`)
4. In AWS Console: **Secrets Manager → `seahaven/notion/api-key` → Edit** — set `apiKey` to your token
5. Run a manual test:
```bash
aws lambda invoke \
--function-name seahaven-notion-sync \
--region us-east-1 \
--log-type Tail \
--query 'LogResult' \
--output text \
/dev/null | base64 -d
```
The sync runs automatically every day at 02:00 UTC via EventBridge.
3. Copy the integration token and update: **Secrets Manager → `seahaven/notion/api-key` → Edit**
### 4. Configure Slack app
In the [Slack API dashboard](https://api.slack.com/apps):
1. **Event Subscriptions** → enable → set Request URL to `https://bot.seahaven.com/slack/events`
2. Subscribe to bot event: `message.im`
3. **OAuth & Permissions** → Bot Token Scopes: `chat:write`, `im:history`
4. **App Home** → Messages Tab → enable, allow users to send messages
5. Reinstall to workspace if prompted
1. **Basic Information** → update Display Name to "Alex", upload avatar
2. **Socket Mode** → enable, generate app-level token with `connections:write` scope
3. Store the app-level token in Secrets Manager as `seahaven/slack/app-level-token`
4. **Event Subscriptions** → enable (no Request URL needed with Socket Mode)
5. Subscribe to bot events: `message.im`, `app_mention`, `app_home_opened`
6. **OAuth & Permissions** → Bot Token Scopes: `chat:write`, `im:history`, `app_mentions:read`, `channels:history`, `groups:history`
7. **App Home** → enable Home Tab and Messages Tab
8. Reinstall app to workspace
9. Invite Alex to channels: `/invite @Alex`
## Project Structure
@ -150,39 +151,38 @@ In the [Slack API dashboard](https://api.slack.com/apps):
bin/
seahaven-slack-bot.ts CDK app entry point
lib/
seahaven-slack-bot-stack.ts Main stack
seahaven-slack-bot-stack.ts Main stack
constructs/
knowledge-base.ts Bedrock KB + AOSS + S3 (via @cdklabs L2 construct)
bedrock-agent.ts Bedrock Agent + QBO/Maps/WO-PO-Site action groups
slack-handler.ts API Gateway + webhook/processor Lambdas + QBO OAuth Lambda
conversation-log.ts DynamoDB table
notion-sync.ts EventBridge daily cron + notion-sync Lambda + Secrets Manager
po-sync.ts EventBridge daily cron + po-sync Lambda
workorder-sync.ts EventBridge daily cron + workorder-sync Lambda
knowledge-base.ts Bedrock KB + AOSS + S3 (via @cdklabs L2 construct)
bedrock-agent.ts Alex — Bedrock Agent + QBO/Maps/WO-PO-Site-Payment action groups
slack-handler.ts Processor Lambda + App Home Lambda + QBO OAuth + API Gateway
socket-mode.ts ECS Fargate service running Slack Socket Mode client
conversation-log.ts DynamoDB conversation + unanswered questions tables
notion-sync.ts EventBridge daily cron + notion-sync Lambda
po-sync.ts EventBridge daily cron + po-sync Lambda
workorder-sync.ts EventBridge daily cron + workorder-sync Lambda
services/
socket-mode/ Socket Mode ECS service (Dockerfile, TypeScript, @slack/socket-mode)
lambda/
slack-webhook/ Verifies Slack signature, fires processor async
slack-processor/ Calls agent, writes DynamoDB, posts Slack reply
qbo-oauth/ OAuth 2.0 connect/callback/disconnect/launch for QuickBooks
qbo-lookup/ Bedrock action group — QuickBooks vendor search
maps-lookup/ Bedrock action group — Google Maps Places search
wo-po-lookup/ Bedrock action group — WO, PO, and site code lookups (DynamoDB direct)
notion-sync/ Daily Notion → S3 export, triggers KB ingestion
po-sync/ Daily purchase-orders DynamoDB → S3 export, triggers KB ingestion
workorder-sync/ Daily WorkOrders DynamoDB → S3 export, triggers KB ingestion
slack-processor/ Calls agent, writes DynamoDB, posts Slack reply
app-home/ Publishes Block Kit capabilities view to App Home tab
qbo-oauth/ OAuth 2.0 connect/callback/disconnect/launch for QuickBooks
qbo-lookup/ Bedrock action group — QuickBooks vendor search
maps-lookup/ Bedrock action group — Google Maps Places search
wo-po-lookup/ Bedrock action group — WO, PO, site, and payment lookups (DynamoDB direct)
notion-sync/ Daily Notion → S3 export, triggers KB ingestion
po-sync/ Daily purchase-orders DynamoDB → S3 export, triggers KB ingestion
workorder-sync/ Daily WorkOrders DynamoDB → S3 export, triggers KB ingestion
scripts/
create-aoss-index.ts Manual fallback for AOSS index creation (not needed in normal deploy)
seed-sites.ts Load Amazon site assignments CSV into DynamoDB
create-aoss-index.ts Manual fallback for AOSS index creation (not needed in normal deploy)
```
## Known Maintenance Items
- **QBO OAuth** — the refresh token auto-rotates on every API call (persisted back to Secrets Manager). If the token ever expires (100 days of inactivity), reconnect via `https://bot.seahaven.com/qbo/connect`. To disconnect, visit `/qbo/disconnect`.
- **Notion sync** runs daily at 02:00 UTC automatically. To trigger an immediate sync, invoke `seahaven-notion-sync` manually via the Lambda console or CLI.
- **PO sync** runs daily at 02:00 UTC. Scans the `purchase-orders` DynamoDB table (from [po-ingest](https://github.com/Sea-Haven-Industries/po-ingest)) and exports each PO as markdown to the KB. Manual trigger: invoke `seahaven-po-sync`.
- **Work order sync** runs daily at 02:00 UTC. Scans `WorkOrders` and `WorkOrderComments` DynamoDB tables (from [workorder-ingest](https://github.com/Sea-Haven-Industries/workorder-ingest)) and exports each work order + comment history as markdown to the KB. Manual trigger: invoke `seahaven-workorder-sync`.
- **Site assignments** are stored in the `SiteAssignments` DynamoDB table. To update, re-run the seed script with the latest CSV:
```bash
npx tsx scripts/seed-sites.ts path/to/updated-site-list.csv
```
The site markdown file in the KB bucket (`sites/amazon-site-assignments.md`) should also be regenerated and re-uploaded, then trigger a KB ingestion job.
- **KB sync for manual S3 uploads** (non-Notion docs) must still be triggered manually after uploading new documents to the S3 bucket.
- **QBO OAuth** — the refresh token auto-rotates on every API call (persisted back to Secrets Manager). If the token ever expires (100 days of inactivity), reconnect via `https://bot.seahaven.com/qbo/connect`.
- **Notion sync** runs daily at 02:00 UTC automatically. Manual trigger: invoke `seahaven-notion-sync`.
- **PO sync** runs daily at 02:00 UTC. Scans `purchase-orders` DynamoDB table (from [po-ingest](https://github.com/Sea-Haven-Industries/po-ingest)). Manual trigger: invoke `seahaven-po-sync`.
- **Work order sync** runs daily at 02:00 UTC. Scans `WorkOrders` and `WorkOrderComments` (from [workorder-ingest](https://github.com/Sea-Haven-Industries/workorder-ingest)). Manual trigger: invoke `seahaven-workorder-sync`.
- **Site assignments** are auto-populated from the `verified-sites` DynamoDB table, maintained by the po-ingest DynamoDB Streams pipeline. No manual seeding required.
- **Unanswered questions** are logged to `seahaven-unanswered-questions` when Alex detects it couldn't answer a question. Review via DynamoDB console — records have `status: "pending"` and 180-day TTL.
- **Socket Mode** — the ECS Fargate task maintains a persistent WebSocket connection to Slack. Monitor via CloudWatch Logs (`socket-mode` log stream). The service auto-restarts on failure.

123
lambda/app-home/index.ts Normal file
View file

@ -0,0 +1,123 @@
import { SecretsManagerClient, GetSecretValueCommand } from '@aws-sdk/client-secrets-manager';
const secretsClient = new SecretsManagerClient({});
interface SlackCredentials {
botToken: string;
signingSecret: string;
}
let cachedCredentials: SlackCredentials | undefined;
async function getSlackCredentials(): Promise<SlackCredentials> {
if (!cachedCredentials) {
const res = await secretsClient.send(
new GetSecretValueCommand({ SecretId: process.env.SLACK_SECRET_ARN! }),
);
cachedCredentials = JSON.parse(res.SecretString!) as SlackCredentials;
}
return cachedCredentials;
}
interface AppHomeEvent {
userId: string;
}
export const handler = async (event: AppHomeEvent): Promise<void> => {
const credentials = await getSlackCredentials();
const blocks = [
{
type: 'header',
text: { type: 'plain_text', text: "Hi, I'm Alex!", emoji: true },
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: "I'm Sea Haven's internal AI assistant. Here's what I can help you with:",
},
},
{ type: 'divider' },
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Vendor Lookup*\nFind existing vendors in QuickBooks or discover new ones nearby.\n_Try: "Find me a plumber near ABE2"_',
},
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Work Order Status*\nLook up any work order by its WO number.\n_Try: "What\'s the status of WO 10046966057?"_',
},
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Purchase Order Details*\nLook up PO status, line items, and shipping info.\n_Try: "Show me PO 2D-20023475"_',
},
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Payment & Invoice Status*\nCheck if a vendor has been paid or if an invoice is scheduled.\n_Try: "Has invoice 12345 been paid?" or "What payments went to Acme?"_',
},
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Site Lookups*\nFind the address and details for any Amazon facility by site code.\n_Try: "Where is DFW6?" or "List all sites in Texas"_',
},
},
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*Company Policies & SOPs*\nAnswer questions about the employee handbook, SA8000 compliance, and SOPs.\n_Try: "What\'s our PTO policy?"_',
},
},
{ type: 'divider' },
{
type: 'section',
text: {
type: 'mrkdwn',
text: '*How to reach me:*\n• Send me a DM right here\n• @mention me in any channel I\'m in',
},
},
{
type: 'context',
elements: [
{
type: 'mrkdwn',
text: "I only know about data in our systems. If I can't answer something, I'll let you know — and my team reviews those gaps regularly to keep improving.",
},
],
},
];
const res = await fetch('https://slack.com/api/views.publish', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${credentials.botToken}`,
},
body: JSON.stringify({
user_id: event.userId,
view: { type: 'home', blocks },
}),
});
if (!res.ok) {
throw new Error(`Slack views.publish HTTP error: ${res.status}`);
}
const data = (await res.json()) as { ok: boolean; error?: string };
if (!data.ok) {
throw new Error(`Slack views.publish API error: ${data.error}`);
}
};

View file

@ -21,6 +21,7 @@ interface ProcessorEvent {
text: string;
ts: string;
threadTs?: string;
eventType?: string;
}
// Cache the secret across warm invocations
@ -142,20 +143,44 @@ function toSlackMrkdwn(md: string): string {
.replace(/\n{3,}/g, '\n\n');
}
const UNANSWERED_SIGNALS = [
"i don't have",
"not in the knowledge base",
"i couldn't find",
"i could not find",
"no information available",
"outside my capabilities",
"no results found",
"i'm unable to",
"i am unable to",
"i don't have access",
"not available in",
"not something i can",
"beyond what i can",
"i don't currently have",
];
function isUnanswered(response: string): boolean {
const lower = response.toLowerCase();
return UNANSWERED_SIGNALS.some((signal) => lower.includes(signal));
}
export const handler = async (event: ProcessorEvent): Promise<void> => {
const { userId, channelId, text, ts, threadTs } = event;
const { userId, channelId, text, ts, threadTs, eventType } = event;
const [credentials, sessionId] = await Promise.all([
getSlackCredentials(),
resolveSessionId(userId),
]);
// Post placeholder immediately so the user sees a response while the agent thinks
// In channels, reply in a thread; in DMs, reply flat (no thread_ts)
const replyThreadTs = eventType === 'app_mention' ? (threadTs ?? ts) : threadTs;
const placeholderTs = await postToSlack(
credentials.botToken,
channelId,
'_Sea Haven Assistant is thinking..._',
threadTs ?? ts,
'_Alex is thinking..._',
replyThreadTs,
);
let answer: string;
@ -182,8 +207,29 @@ export const handler = async (event: ProcessorEvent): Promise<void> => {
answer,
channelId,
messageTs: ts,
eventType: eventType ?? 'dm',
ttl,
},
}),
);
if (process.env.UNANSWERED_TABLE && isUnanswered(answer)) {
const qTtl = Math.floor(Date.now() / 1000) + 180 * 24 * 60 * 60;
await ddb.send(
new PutCommand({
TableName: process.env.UNANSWERED_TABLE,
Item: {
questionId: crypto.randomUUID(),
userId,
question: text,
agentResponse: answer,
timestamp: now,
status: 'pending',
channelId,
eventType: eventType ?? 'dm',
ttl: qTtl,
},
}),
);
}
};

View file

@ -1,115 +0,0 @@
import type { APIGatewayProxyEventV2, APIGatewayProxyResultV2 } from 'aws-lambda';
import { LambdaClient, InvokeCommand } from '@aws-sdk/client-lambda';
import { SecretsManagerClient, GetSecretValueCommand } from '@aws-sdk/client-secrets-manager';
import { createHmac, timingSafeEqual } from 'crypto';
const lambdaClient = new LambdaClient({});
const secretsClient = new SecretsManagerClient({});
interface SlackCredentials {
botToken: string;
signingSecret: string;
}
// Cache the secret within the Lambda execution environment lifetime
let cachedSecret: SlackCredentials | undefined;
async function getSlackCredentials(): Promise<SlackCredentials> {
if (!cachedSecret) {
const res = await secretsClient.send(
new GetSecretValueCommand({ SecretId: process.env.SLACK_SECRET_ARN! }),
);
cachedSecret = JSON.parse(res.SecretString!) as SlackCredentials;
}
return cachedSecret;
}
function verifySignature(
signingSecret: string,
timestamp: string,
rawBody: string,
signature: string,
): boolean {
// Reject requests older than 5 minutes (replay attack prevention)
const ageSeconds = Math.abs(Date.now() / 1000 - parseInt(timestamp, 10));
if (ageSeconds > 300) return false;
const baseString = `v0:${timestamp}:${rawBody}`;
const expected = `v0=${createHmac('sha256', signingSecret).update(baseString).digest('hex')}`;
try {
return timingSafeEqual(Buffer.from(expected, 'utf8'), Buffer.from(signature, 'utf8'));
} catch {
return false;
}
}
export const handler = async (
event: APIGatewayProxyEventV2,
): Promise<APIGatewayProxyResultV2> => {
const rawBody = event.body ?? '';
const timestamp = event.headers['x-slack-request-timestamp'] ?? '';
const signature = event.headers['x-slack-signature'] ?? '';
const credentials = await getSlackCredentials();
if (!verifySignature(credentials.signingSecret, timestamp, rawBody, signature)) {
return { statusCode: 401, body: 'Invalid signature' };
}
const payload = JSON.parse(rawBody) as Record<string, unknown>;
// Slack sends this when you first register the event URL
if (payload.type === 'url_verification') {
return {
statusCode: 200,
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ challenge: payload.challenge }),
};
}
if (payload.type !== 'event_callback') {
return { statusCode: 200, body: 'OK' };
}
const slackEvent = payload.event as Record<string, unknown>;
// DM only — channel_type === 'im'
if (slackEvent.channel_type !== 'im') {
return { statusCode: 200, body: 'OK' };
}
// Ignore bot messages and message edits/deletions to prevent loops
if (
slackEvent.bot_id ||
slackEvent.subtype === 'bot_message' ||
slackEvent.subtype === 'message_changed' ||
slackEvent.subtype === 'message_deleted'
) {
return { statusCode: 200, body: 'OK' };
}
if (slackEvent.type !== 'message') {
return { statusCode: 200, body: 'OK' };
}
// Fire-and-forget — processor handles the slow Bedrock call
await lambdaClient.send(
new InvokeCommand({
FunctionName: process.env.PROCESSOR_FUNCTION_NAME!,
InvocationType: 'Event',
Payload: Buffer.from(
JSON.stringify({
userId: slackEvent.user,
channelId: slackEvent.channel,
text: slackEvent.text,
ts: slackEvent.ts,
threadTs: slackEvent.thread_ts,
}),
),
}),
);
// Slack requires a 200 within 3 seconds
return { statusCode: 200, body: 'OK' };
};

View file

@ -1,4 +1,4 @@
import { DynamoDBClient, GetItemCommand, QueryCommand } from '@aws-sdk/client-dynamodb';
import { DynamoDBClient, GetItemCommand, QueryCommand, ScanCommand } from '@aws-sdk/client-dynamodb';
import { unmarshall } from '@aws-sdk/util-dynamodb';
const dynamo = new DynamoDBClient({});
@ -248,6 +248,94 @@ async function lookupSite(siteCode: string, state?: string): Promise<string> {
return 'Please provide a site code (e.g., "ABE2") or a state abbreviation (e.g., "TX") to search.';
}
// ── Payment lookups ─────────────────────────────────────────────────────────
function formatPayment(p: Record<string, any>): string {
const lines: string[] = [];
lines.push(`Check #${p.check_number ?? 'N/A'} — ${p.payee ?? 'Unknown payee'}`);
if (p.amount_usd !== undefined) lines.push(`Amount: $${p.amount_usd}`);
if (p.method) lines.push(`Method: ${p.method}`);
if (p.status) lines.push(`Status: ${p.status}`);
if (p.send_payment_on) lines.push(`Send Date: ${p.send_payment_on}`);
if (p.clear_status) lines.push(`Clear Status: ${p.clear_status}`);
if (p.cleared_date) lines.push(`Cleared Date: ${p.cleared_date}`);
if (p.invoice_numbers) lines.push(`Invoices: ${p.invoice_numbers}`);
if (p.company_subsidiary) lines.push(`Subsidiary: ${p.company_subsidiary}`);
if (p.bank_reference) lines.push(`Bank Ref: ${p.bank_reference}`);
return lines.join('\n');
}
async function scanPayments(): Promise<Record<string, any>[]> {
const items: Record<string, any>[] = [];
let lastKey: Record<string, any> | undefined;
do {
const res = await dynamo.send(
new ScanCommand({
TableName: process.env.PAYMENTS_TABLE!,
FilterExpression: 'begins_with(pk, :prefix)',
ExpressionAttributeValues: { ':prefix': { S: 'payment#' } },
ExclusiveStartKey: lastKey,
}),
);
for (const item of res.Items ?? []) {
items.push(unmarshall(item));
}
lastKey = res.LastEvaluatedKey;
} while (lastKey);
return items;
}
async function lookupPaymentByVendor(vendorName: string): Promise<string> {
const all = await scanPayments();
const needle = vendorName.toLowerCase();
const matches = all.filter((p) => (p.payee ?? '').toLowerCase().includes(needle));
if (matches.length === 0) {
return `No payments found for vendor matching "${vendorName}" in the last 90 days.`;
}
const lines = [`Payments matching "${vendorName}" (${matches.length} found):\n`];
for (const p of matches) {
lines.push(formatPayment(p));
lines.push('');
}
return lines.join('\n');
}
async function lookupPaymentByInvoice(invoiceNumber: string): Promise<string> {
const all = await scanPayments();
const needle = invoiceNumber.toLowerCase();
const matches = all.filter((p) => (p.invoice_numbers ?? '').toLowerCase().includes(needle));
if (matches.length === 0) {
return `No payments found for invoice "${invoiceNumber}" in the last 90 days.`;
}
const lines = [`Payments for invoice "${invoiceNumber}" (${matches.length} found):\n`];
for (const p of matches) {
lines.push(formatPayment(p));
lines.push('');
}
return lines.join('\n');
}
async function lookupPaymentByCheck(checkNumber: string): Promise<string> {
const res = await dynamo.send(
new GetItemCommand({
TableName: process.env.PAYMENTS_TABLE!,
Key: { pk: { S: `payment#${checkNumber}` } },
}),
);
if (!res.Item) {
return `No payment found for check number "${checkNumber}".`;
}
return formatPayment(unmarshall(res.Item));
}
// ── Handler ──────────────────────────────────────────────────────────────────
export const handler = async (event: BedrockActionEvent): Promise<BedrockActionResponse> => {
@ -271,6 +359,18 @@ export const handler = async (event: BedrockActionEvent): Promise<BedrockActionR
const state = params.state?.trim();
if (!code && !state) return makeResponse(event, 'Please provide a site code or state abbreviation.');
result = await lookupSite(code, state);
} else if (event.function === 'lookup_payment_by_vendor') {
const name = params.vendor_name?.trim();
if (!name) return makeResponse(event, 'Please provide a vendor or payee name.');
result = await lookupPaymentByVendor(name);
} else if (event.function === 'lookup_payment_by_invoice') {
const inv = params.invoice_number?.trim();
if (!inv) return makeResponse(event, 'Please provide an invoice number.');
result = await lookupPaymentByInvoice(inv);
} else if (event.function === 'lookup_payment_by_check') {
const check = params.check_number?.trim();
if (!check) return makeResponse(event, 'Please provide a check number.');
result = await lookupPaymentByCheck(check);
} else {
result = `Unknown function: ${event.function}`;
}

View file

@ -3,6 +3,7 @@ import { Construct } from 'constructs';
import * as iam from 'aws-cdk-lib/aws-iam';
import * as lambda from 'aws-cdk-lib/aws-lambda';
import * as lambdaNodejs from 'aws-cdk-lib/aws-lambda-nodejs';
import * as logs from 'aws-cdk-lib/aws-logs';
import * as dynamodb from 'aws-cdk-lib/aws-dynamodb';
import * as ec2 from 'aws-cdk-lib/aws-ec2';
import * as secretsmanager from 'aws-cdk-lib/aws-secretsmanager';
@ -51,6 +52,8 @@ export class BedrockAgentConstruct extends Construct {
entry: path.join(__dirname, '../../lambda/qbo-lookup/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.seconds(30),
memorySize: 256,
environment: { QBO_SECRET_ARN: qboSecret.secretArn },
@ -68,6 +71,8 @@ export class BedrockAgentConstruct extends Construct {
entry: path.join(__dirname, '../../lambda/maps-lookup/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.seconds(30),
memorySize: 256,
environment: { MAPS_SECRET_ARN: mapsSecret.secretArn },
@ -86,24 +91,21 @@ export class BedrockAgentConstruct extends Construct {
this, 'PurchaseOrdersTable', 'purchase-orders',
);
// Site assignments table — stores Amazon facility site codes and addresses
const sitesTable = new dynamodb.Table(this, 'SiteAssignmentsTable', {
tableName: 'SiteAssignments',
partitionKey: { name: 'siteCode', type: dynamodb.AttributeType.STRING },
billingMode: dynamodb.BillingMode.PAY_PER_REQUEST,
removalPolicy: cdk.RemovalPolicy.RETAIN,
});
sitesTable.addGlobalSecondaryIndex({
indexName: 'by-state',
partitionKey: { name: 'state', type: dynamodb.AttributeType.STRING },
projectionType: dynamodb.ProjectionType.ALL,
});
const sitesTable = dynamodb.Table.fromTableName(
this, 'VerifiedSitesTable', 'verified-sites',
);
const paymentsTable = dynamodb.Table.fromTableName(
this, 'PaymentsDashboardTable', 'PaymentsDashboard',
);
this.woPoLambda = new lambdaNodejs.NodejsFunction(this, 'WoPoLookupFn', {
functionName: 'seahaven-wo-po-lookup',
entry: path.join(__dirname, '../../lambda/wo-po-lookup/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.seconds(30),
memorySize: 256,
environment: {
@ -111,6 +113,7 @@ export class BedrockAgentConstruct extends Construct {
COMMENTS_TABLE: commentsTable.tableName,
PO_TABLE: poTable.tableName,
SITES_TABLE: sitesTable.tableName,
PAYMENTS_TABLE: paymentsTable.tableName,
},
bundling,
});
@ -118,6 +121,7 @@ export class BedrockAgentConstruct extends Construct {
commentsTable.grantReadData(this.woPoLambda);
poTable.grantReadData(this.woPoLambda);
sitesTable.grantReadData(this.woPoLambda);
paymentsTable.grantReadData(this.woPoLambda);
// ── Bedrock Agent execution role ──────────────────────────────────────────
const agentRole = new iam.Role(this, 'AgentRole', {
@ -161,7 +165,7 @@ export class BedrockAgentConstruct extends Construct {
});
// ── Agent instruction (system prompt) ────────────────────────────────────
const instruction = `You are the Sea Haven Industries internal assistant, accessible to employees via Slack direct message. Sea Haven is a facility services company that provides maintenance, repair, and facility services to clients like Amazon.
const instruction = `Your name is Alex. You are Sea Haven Industries' internal AI assistant, available to employees via Slack DM or @mention in any channel. You are friendly, concise, and professional — like a helpful coworker, not a robot. Sea Haven is a facility services company that provides maintenance, repair, and facility services to clients like Amazon.
IMPORTANT: You work FOR Sea Haven. Sea Haven is our company — we are the vendor/contractor that gets dispatched to job sites. When work orders or comments mention "Sea Haven" being dispatched or assigned, that means OUR team was sent. Never suggest Sea Haven as an external vendor to contact — we ARE Sea Haven. When users ask for a "local vendor," they mean a subcontractor or specialty trade vendor to handle work on our behalf.
@ -172,6 +176,7 @@ You help employees with:
4. Work order lookups — status, history, assignments, comments, and details for any work order by its WO number
5. Purchase order lookups — status, line items, suppliers, ship-to details, and dates for any PO by its PO number
6. Amazon site lookups — address, location, and details for any Amazon facility by its site code (e.g., "ABE2", "DFW6"), or listing all sites in a given state
7. Payment and invoice status — whether a vendor has been paid, if an invoice is scheduled, check details
## Vendor Query Rules — STRICTLY follow this priority order:
@ -189,6 +194,9 @@ When presenting vendor results:
## Work Order & Purchase Order Queries:
When a user asks about a work order (WO) or purchase order (PO) by number, ALWAYS use the WO_PO_Lookup action group to retrieve the record directly. Do NOT use the knowledge base for WO/PO lookups by number — the action group queries the database directly and is more reliable. Present the returned data clearly and concisely. The output is already well-structured — relay the key information without adding excessive formatting or repeating section headers verbatim. Summarize the current status and most recent updates first, then include the full comment history.
## Payment & Invoice Queries:
When a user asks about payments, invoices, or whether a vendor has been paid, use the WO_PO_Lookup payment functions. For vendor payment searches, use lookup_payment_by_vendor. For invoice number lookups, use lookup_payment_by_invoice. For check number lookups, use lookup_payment_by_check. Present results clearly: check number, payee, amount, status, method, and relevant dates. Payment records cover the last 90 days.
## Site Lookups:
When a user asks about an Amazon site by its code (e.g., "ABE2", "DFW6", "WWY1"), ALWAYS use the WO_PO_Lookup.lookup_site function. You can also list all sites in a state by passing the state abbreviation. Do NOT use the knowledge base for site code lookups — the action group queries the database directly and is more reliable.
@ -199,8 +207,8 @@ Keep responses concise, professional, and actionable.`;
// ── CfnAgent ──────────────────────────────────────────────────────────────
this.agent = new bedrock.CfnAgent(this, 'Agent', {
agentName: 'seahaven-assistant',
description: 'Sea Haven Industries internal Slack assistant',
agentName: 'seahaven-alex',
description: 'Alex — Sea Haven Industries internal Slack assistant',
agentResourceRoleArn: agentRole.roleArn,
foundationModel: BedrockAgentConstruct.MODEL_ID,
instruction,
@ -267,7 +275,7 @@ Keep responses concise, professional, and actionable.`;
},
{
actionGroupName: 'WO_PO_Lookup',
description: 'Look up work orders, purchase orders, and Amazon site assignments directly from the database',
description: 'Look up work orders, purchase orders, Amazon site assignments, and payment/invoice status directly from the database',
actionGroupState: 'ENABLED',
actionGroupExecutor: { lambda: this.woPoLambda.functionArn },
functionSchema: {
@ -310,6 +318,39 @@ Keep responses concise, professional, and actionable.`;
},
},
},
{
name: 'lookup_payment_by_vendor',
description: 'Search for payments made to a vendor/payee by name. Returns check numbers, amounts, status, payment method, and dates. Use when an employee asks "Has vendor X been paid?" or "What payments went to X?"',
parameters: {
vendor_name: {
type: 'string',
description: 'The vendor or payee name to search for (partial match supported)',
required: true,
},
},
},
{
name: 'lookup_payment_by_invoice',
description: 'Search for a payment by invoice number. Returns the check number, payee, amount, status, and payment date. Use when an employee asks "Is invoice 12345 scheduled?" or "Has invoice 12345 been paid?"',
parameters: {
invoice_number: {
type: 'string',
description: 'The invoice number to search for',
required: true,
},
},
},
{
name: 'lookup_payment_by_check',
description: 'Look up a specific payment by its check number. Returns payee, amount, status, method, invoices covered, and dates.',
parameters: {
check_number: {
type: 'string',
description: 'The check number to look up',
required: true,
},
},
},
],
},
},
@ -321,7 +362,7 @@ Keep responses concise, professional, and actionable.`;
this.agentAlias = new bedrock.CfnAgentAlias(this, 'AgentAlias', {
agentId: this.agent.attrAgentId,
agentAliasName: 'live',
description: 'Production alias — seahaven-assistant v8 (site lookups)',
description: 'Production alias — Alex v1',
});
}
}

View file

@ -4,6 +4,7 @@ import * as dynamodb from 'aws-cdk-lib/aws-dynamodb';
export class ConversationLogConstruct extends Construct {
public readonly table: dynamodb.Table;
public readonly unansweredTable: dynamodb.Table;
constructor(scope: Construct, id: string) {
super(scope, id);
@ -27,5 +28,13 @@ export class ConversationLogConstruct extends Construct {
partitionKey: { name: 'sessionId', type: dynamodb.AttributeType.STRING },
sortKey: { name: 'timestamp', type: dynamodb.AttributeType.STRING },
});
this.unansweredTable = new dynamodb.Table(this, 'UnansweredTable', {
tableName: 'seahaven-unanswered-questions',
partitionKey: { name: 'questionId', type: dynamodb.AttributeType.STRING },
billingMode: dynamodb.BillingMode.PAY_PER_REQUEST,
removalPolicy: cdk.RemovalPolicy.RETAIN,
timeToLiveAttribute: 'ttl',
});
}
}

View file

@ -2,6 +2,7 @@ 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 logs from 'aws-cdk-lib/aws-logs';
import * as s3 from 'aws-cdk-lib/aws-s3';
import * as secretsmanager from 'aws-cdk-lib/aws-secretsmanager';
import * as events from 'aws-cdk-lib/aws-events';
@ -38,6 +39,8 @@ export class NotionSyncConstruct extends Construct {
functionName: 'seahaven-notion-sync',
entry: path.join(__dirname, '../../lambda/notion-sync/index.ts'),
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
memorySize: 512,
timeout: cdk.Duration.minutes(5),
environment: {

View file

@ -2,6 +2,7 @@ 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 logs from 'aws-cdk-lib/aws-logs';
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';
@ -32,6 +33,8 @@ export class PoSyncConstruct extends Construct {
functionName: 'seahaven-po-sync',
entry: path.join(__dirname, '../../lambda/po-sync/index.ts'),
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
memorySize: 512,
timeout: cdk.Duration.minutes(15),
environment: {

View file

@ -2,6 +2,7 @@ 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 logs from 'aws-cdk-lib/aws-logs';
import * as apigatewayv2 from 'aws-cdk-lib/aws-apigatewayv2';
import { HttpLambdaIntegration } from 'aws-cdk-lib/aws-apigatewayv2-integrations';
import * as dynamodb from 'aws-cdk-lib/aws-dynamodb';
@ -19,97 +20,88 @@ export interface SlackHandlerProps {
agentId: string;
agentAliasId: string;
conversationTable: dynamodb.Table;
unansweredTable: dynamodb.Table;
wildcardCertArn: string;
vpc: ec2.IVpc;
lambdaSecurityGroup: ec2.ISecurityGroup;
}
export class SlackHandlerConstruct extends Construct {
public readonly webhookLambda: lambdaNodejs.NodejsFunction;
public readonly processorLambda: lambdaNodejs.NodejsFunction;
public readonly appHomeLambda: lambdaNodejs.NodejsFunction;
public readonly api: apigatewayv2.HttpApi;
constructor(scope: Construct, id: string, props: SlackHandlerProps) {
super(scope, id);
// Slack credentials secret (created manually — see README for structure)
const slackSecret = secretsmanager.Secret.fromSecretNameV2(
this, 'SlackSecret', 'seahaven/slack/credentials',
);
const bundling: lambdaNodejs.BundlingOptions = {
externalModules: ['@aws-sdk/*'],
minify: true,
sourceMap: false,
};
// ── Processor Lambda ──────────────────────────────────────────────────────
// Async worker: calls Bedrock Agent, writes to DynamoDB, posts reply to Slack.
// Timeout is generous — agent invocations with multi-step tool use can take 2–3 min.
this.processorLambda = new lambdaNodejs.NodejsFunction(this, 'ProcessorFn', {
functionName: 'seahaven-slack-processor',
entry: path.join(__dirname, '../../lambda/slack-processor/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.minutes(5),
memorySize: 512,
environment: {
AGENT_ID: props.agentId,
AGENT_ALIAS_ID: props.agentAliasId,
CONVERSATION_TABLE: props.conversationTable.tableName,
UNANSWERED_TABLE: props.unansweredTable.tableName,
SLACK_SECRET_ARN: slackSecret.secretArn,
REGION: props.region,
},
bundling: {
externalModules: ['@aws-sdk/*'],
minify: true,
sourceMap: false,
},
bundling,
});
slackSecret.grantRead(this.processorLambda);
props.conversationTable.grantReadWriteData(this.processorLambda);
props.unansweredTable.grantWriteData(this.processorLambda);
this.processorLambda.addToRolePolicy(new iam.PolicyStatement({
actions: ['bedrock:InvokeAgent'],
resources: [
// Wildcard on agent alias — agentAliasId is a CDK token resolved at synth
`arn:aws:bedrock:${props.region}:${props.accountId}:agent/${props.agentId}`,
`arn:aws:bedrock:${props.region}:${props.accountId}:agent-alias/${props.agentId}/*`,
],
}));
// ── Webhook Lambda ────────────────────────────────────────────────────────
// Synchronous: verifies Slack signature, returns 200 immediately, fires processor async.
this.webhookLambda = new lambdaNodejs.NodejsFunction(this, 'WebhookFn', {
functionName: 'seahaven-slack-webhook',
entry: path.join(__dirname, '../../lambda/slack-webhook/index.ts'),
// ── App Home Lambda ───────────────────────────────────────────────────────
this.appHomeLambda = new lambdaNodejs.NodejsFunction(this, 'AppHomeFn', {
functionName: 'seahaven-app-home',
entry: path.join(__dirname, '../../lambda/app-home/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.seconds(10),
memorySize: 256,
environment: {
PROCESSOR_FUNCTION_NAME: this.processorLambda.functionName,
SLACK_SECRET_ARN: slackSecret.secretArn,
},
bundling: {
externalModules: ['@aws-sdk/*'],
minify: true,
sourceMap: false,
},
bundling,
});
slackSecret.grantRead(this.webhookLambda);
this.processorLambda.grantInvoke(this.webhookLambda);
slackSecret.grantRead(this.appHomeLambda);
// ── HTTP API (API Gateway v2) ──────────────────────────────────────────────
// ── HTTP API (API Gateway v2) — QBO OAuth routes only ─────────────────────
this.api = new apigatewayv2.HttpApi(this, 'Api', {
apiName: 'seahaven-slack-webhook',
description: 'Receives Slack event webhook calls for Sea Haven bot',
});
this.api.addRoutes({
path: '/slack/events',
methods: [apigatewayv2.HttpMethod.POST],
integration: new HttpLambdaIntegration('WebhookIntegration', this.webhookLambda),
description: 'Sea Haven bot — QBO OAuth routes',
});
// ── QBO OAuth Lambda ─────────────────────────────────────────────────────
// Handles /qbo/connect, /qbo/callback, /qbo/disconnect, /qbo/launch
const qboSecret = secretsmanager.Secret.fromSecretNameV2(
this, 'QBOSecret', 'seahaven/qbo/oauth',
);
@ -119,6 +111,8 @@ export class SlackHandlerConstruct extends Construct {
entry: path.join(__dirname, '../../lambda/qbo-oauth/index.ts'),
handler: 'handler',
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
timeout: cdk.Duration.seconds(15),
memorySize: 256,
environment: {
@ -128,11 +122,7 @@ export class SlackHandlerConstruct extends Construct {
vpc: props.vpc,
vpcSubnets: { subnetType: ec2.SubnetType.PRIVATE_WITH_EGRESS },
securityGroups: [props.lambdaSecurityGroup],
bundling: {
externalModules: ['@aws-sdk/*'],
minify: true,
sourceMap: false,
},
bundling,
});
qboSecret.grantRead(qboOAuthLambda);
@ -178,10 +168,5 @@ export class SlackHandlerConstruct extends Construct {
),
),
});
new cdk.CfnOutput(scope, 'SlackWebhookUrl', {
value: 'https://bot.seahaven.com/slack/events',
description: 'Paste this into Slack app → Event Subscriptions → Request URL',
});
}
}

View file

@ -0,0 +1,67 @@
import { Construct } from 'constructs';
import * as ecs from 'aws-cdk-lib/aws-ecs';
import * as ec2 from 'aws-cdk-lib/aws-ec2';
import * as secretsmanager from 'aws-cdk-lib/aws-secretsmanager';
import * as lambda from 'aws-cdk-lib/aws-lambda';
import * as logs from 'aws-cdk-lib/aws-logs';
import * as path from 'path';
export interface SocketModeProps {
vpc: ec2.IVpc;
processorLambda: lambda.IFunction;
appHomeLambda: lambda.IFunction;
}
export class SocketModeConstruct extends Construct {
constructor(scope: Construct, id: string, props: SocketModeProps) {
super(scope, id);
const appTokenSecret = secretsmanager.Secret.fromSecretNameV2(
this, 'AppTokenSecret', 'seahaven/slack/app-level-token',
);
const cluster = new ecs.Cluster(this, 'Cluster', {
clusterName: 'seahaven-socket-mode',
vpc: props.vpc,
});
const taskDef = new ecs.FargateTaskDefinition(this, 'TaskDef', {
memoryLimitMiB: 512,
cpu: 256,
runtimePlatform: {
cpuArchitecture: ecs.CpuArchitecture.ARM64,
operatingSystemFamily: ecs.OperatingSystemFamily.LINUX,
},
});
taskDef.addContainer('socket-mode', {
image: ecs.ContainerImage.fromAsset(
path.join(__dirname, '../../services/socket-mode'),
),
environment: {
APP_TOKEN_SECRET_ARN: appTokenSecret.secretArn,
PROCESSOR_FUNCTION_NAME: props.processorLambda.functionName,
APP_HOME_FUNCTION_NAME: props.appHomeLambda.functionName,
},
logging: ecs.LogDrivers.awsLogs({
streamPrefix: 'socket-mode',
logRetention: logs.RetentionDays.TWO_MONTHS,
}),
});
appTokenSecret.grantRead(taskDef.taskRole);
props.processorLambda.grantInvoke(taskDef.taskRole);
props.appHomeLambda.grantInvoke(taskDef.taskRole);
new ecs.FargateService(this, 'Service', {
serviceName: 'seahaven-socket-mode',
cluster,
taskDefinition: taskDef,
desiredCount: 1,
minHealthyPercent: 100,
maxHealthyPercent: 200,
assignPublicIp: false,
vpcSubnets: { subnetType: ec2.SubnetType.PRIVATE_WITH_EGRESS },
});
}
}

View file

@ -2,6 +2,7 @@ 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 logs from 'aws-cdk-lib/aws-logs';
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';
@ -35,6 +36,8 @@ export class WorkorderSyncConstruct extends Construct {
functionName: 'seahaven-workorder-sync',
entry: path.join(__dirname, '../../lambda/workorder-sync/index.ts'),
runtime: lambda.Runtime.NODEJS_22_X,
architecture: lambda.Architecture.ARM_64,
logRetention: logs.RetentionDays.TWO_MONTHS,
memorySize: 512,
timeout: cdk.Duration.minutes(5),
environment: {

View file

@ -5,6 +5,7 @@ import { ConversationLogConstruct } from './constructs/conversation-log';
import { KnowledgeBaseConstruct } from './constructs/knowledge-base';
import { BedrockAgentConstruct } from './constructs/bedrock-agent';
import { SlackHandlerConstruct } from './constructs/slack-handler';
import { SocketModeConstruct } from './constructs/socket-mode';
import { NotionSyncConstruct } from './constructs/notion-sync';
import { PoSyncConstruct } from './constructs/po-sync';
import { WorkorderSyncConstruct } from './constructs/workorder-sync';
@ -13,7 +14,6 @@ export class SeahavenSlackBotStack extends cdk.Stack {
constructor(scope: Construct, id: string, props?: cdk.StackProps) {
super(scope, id, props);
// Wildcard cert ARN (*.seahaven.com) — set in cdk.json or via --context
const wildcardCertArn = this.node.tryGetContext('wildcardCertArn') as string | undefined;
if (!wildcardCertArn || wildcardCertArn.includes('CHANGE-ME')) {
throw new Error(
@ -22,7 +22,7 @@ export class SeahavenSlackBotStack extends cdk.Stack {
);
}
// ── VPC (existing) — QBO Lambdas run here for static outbound IP ──────────
// ── VPC (existing) — QBO Lambdas + Socket Mode run here ──────────────────
const vpc = ec2.Vpc.fromLookup(this, 'SeahavenVpc', { vpcId: 'vpc-0d3d4b67bd0cf8a68' });
const lambdaSecurityGroup = new ec2.SecurityGroup(this, 'QBOLambdaSG', {
@ -32,7 +32,7 @@ export class SeahavenSlackBotStack extends cdk.Stack {
allowAllOutbound: true,
});
// ── Conversation history (DynamoDB) ───────────────────────────────────────
// ── Conversation history + unanswered questions (DynamoDB) ────────────────
const conversationLog = new ConversationLogConstruct(this, 'ConversationLog');
// ── Bedrock Knowledge Base (OpenSearch Serverless + S3) ───────────────────
@ -41,7 +41,7 @@ export class SeahavenSlackBotStack extends cdk.Stack {
region: this.region,
});
// ── Bedrock Agent (Claude Sonnet + QBO + Google Maps action groups) ───────
// ── Bedrock Agent (Alex — Claude Sonnet + action groups) ─────────────────
const bedrockAgent = new BedrockAgentConstruct(this, 'BedrockAgent', {
accountId: this.account,
region: this.region,
@ -75,18 +75,26 @@ export class SeahavenSlackBotStack extends cdk.Stack {
dataSourceId: knowledgeBase.dataSource.dataSourceId,
});
// ── Slack webhook handler (API Gateway + Lambda) ───────────────────────────
new SlackHandlerConstruct(this, 'SlackHandler', {
// ── Slack handler (processor + app home + QBO OAuth + API Gateway) ────────
const slackHandler = new SlackHandlerConstruct(this, 'SlackHandler', {
accountId: this.account,
region: this.region,
agentId: bedrockAgent.agent.attrAgentId,
agentAliasId: bedrockAgent.agentAlias.attrAgentAliasId,
conversationTable: conversationLog.table,
unansweredTable: conversationLog.unansweredTable,
wildcardCertArn,
vpc,
lambdaSecurityGroup,
});
// ── Socket Mode (ECS Fargate — replaces webhook Lambda) ──────────────────
new SocketModeConstruct(this, 'SocketMode', {
vpc,
processorLambda: slackHandler.processorLambda,
appHomeLambda: slackHandler.appHomeLambda,
});
// ── Stack outputs ─────────────────────────────────────────────────────────
new cdk.CfnOutput(this, 'KBDocsBucketName', {
value: knowledgeBase.docsBucket.bucketName,
@ -95,8 +103,7 @@ export class SeahavenSlackBotStack extends cdk.Stack {
new cdk.CfnOutput(this, 'AgentId', {
value: bedrockAgent.agent.attrAgentId,
description: 'Bedrock Agent ID',
description: 'Bedrock Agent ID (Alex)',
});
}
}

7
package-lock.json generated
View file

@ -46,6 +46,7 @@
"semver"
],
"license": "Apache-2.0",
"peer": true,
"dependencies": {
"jsonschema": "~1.4.1",
"semver": "^7.7.4"
@ -214,6 +215,7 @@
"integrity": "sha512-mJlCunrAcjOvRyjDiOSNNFEJWwGkfHChqNHZI36oZwnbWyVBkMa43Qhc54sWIhZVXzYONeQ+hviF6zLbFBTUAw==",
"dev": true,
"license": "Apache-2.0",
"peer": true,
"dependencies": {
"@aws-crypto/sha256-browser": "5.2.0",
"@aws-crypto/sha256-js": "5.2.0",
@ -2015,6 +2017,7 @@
"resolved": "https://registry.npmjs.org/@types/node/-/node-22.19.17.tgz",
"integrity": "sha512-wGdMcf+vPYM6jikpS/qhg6WiqSV/OhG+jeeHT/KlVqxYfD40iYJf9/AE1uQxVWFvU7MipKRkRv8NSHiCGgPr8Q==",
"license": "MIT",
"peer": true,
"dependencies": {
"undici-types": "~6.21.0"
}
@ -2520,7 +2523,8 @@
"version": "10.6.0",
"resolved": "https://registry.npmjs.org/constructs/-/constructs-10.6.0.tgz",
"integrity": "sha512-TxHOnBO5zMo/G76ykzGF/wMpEHu257TbWiIxP9K0Yv/+t70UzgBQiTqjkAsWOPC6jW91DzJI0+ehQV6xDRNBuQ==",
"license": "Apache-2.0"
"license": "Apache-2.0",
"peer": true
},
"node_modules/create-require": {
"version": "1.1.1",
@ -3001,6 +3005,7 @@
"integrity": "sha512-84MVSjMEHP+FQRPy3pX9sTVV/INIex71s9TL2Gm5FG/WG1SqXeKyZ0k7/blY/4FdOzI12CBy1vGc4og/eus0fw==",
"dev": true,
"license": "Apache-2.0",
"peer": true,
"bin": {
"tsc": "bin/tsc",
"tsserver": "bin/tsserver"

View file

@ -1,84 +0,0 @@
/**
* Seed the SiteAssignments DynamoDB table from the Amazon site list CSV.
*
* Usage:
* npx tsx scripts/seed-sites.ts path/to/amazon-sites.csv
*/
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, BatchWriteCommand } from '@aws-sdk/lib-dynamodb';
import { readFileSync } from 'fs';
const TABLE_NAME = 'SiteAssignments';
const REGION = 'us-east-1';
const ddb = DynamoDBDocumentClient.from(new DynamoDBClient({ region: REGION }));
function parseCSV(filePath: string) {
const content = readFileSync(filePath, 'utf-8');
const lines = content.split('\n').filter((l) => l.trim());
return lines.slice(1).map((line) => {
const cols = line.split(',');
return {
siteCode: (cols[0] ?? '').trim(),
address: (cols[1] ?? '').trim(),
city: (cols[2] ?? '').trim(),
state: (cols[3] ?? '').trim(),
zip: (cols[4] ?? '').trim(),
fullAddress: (cols[5] ?? '').trim(),
latitude: (cols[6] ?? '').trim(),
longitude: (cols[7] ?? '').trim(),
notes: (cols[9] ?? '').trim(),
};
}).filter((r) => r.siteCode);
}
async function seed(filePath: string) {
const allRows = parseCSV(filePath);
// Deduplicate by siteCode — last occurrence wins
const deduped = new Map(allRows.map((r) => [r.siteCode, r]));
const rows = [...deduped.values()];
console.log(`Parsed ${allRows.length} rows, ${rows.length} unique sites`);
let written = 0;
for (let i = 0; i < rows.length; i += 25) {
const batch = rows.slice(i, i + 25);
const requests = batch.map((row) => ({
PutRequest: {
Item: {
siteCode: row.siteCode,
address: row.address,
city: row.city,
state: row.state,
...(row.zip && { zip: row.zip }),
fullAddress: row.fullAddress,
...(row.latitude && { latitude: row.latitude }),
...(row.longitude && { longitude: row.longitude }),
...(row.notes && { notes: row.notes }),
},
},
}));
await ddb.send(
new BatchWriteCommand({
RequestItems: { [TABLE_NAME]: requests },
}),
);
written += batch.length;
if (written % 250 === 0 || written === rows.length) {
console.log(` Written ${written}/${rows.length}`);
}
}
console.log(`Done — seeded ${written} sites into ${TABLE_NAME}`);
}
const csvPath = process.argv[2];
if (!csvPath) {
console.error('Usage: npx tsx scripts/seed-sites.ts <path-to-csv>');
process.exit(1);
}
seed(csvPath).catch((err) => {
console.error('Seed failed:', err);
process.exit(1);
});

View file

@ -0,0 +1,14 @@
FROM node:22-slim AS builder
WORKDIR /app
COPY package*.json tsconfig.json ./
RUN npm ci
COPY src/ src/
RUN npx tsc
FROM node:22-slim
WORKDIR /app
COPY --from=builder /app/dist ./dist
COPY --from=builder /app/node_modules ./node_modules
COPY package.json ./
USER node
CMD ["node", "dist/index.js"]

1992
services/socket-mode/package-lock.json generated Normal file

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,19 @@
{
"name": "seahaven-socket-mode",
"version": "1.0.0",
"private": true,
"scripts": {
"build": "tsc",
"start": "node dist/index.js"
},
"dependencies": {
"@slack/socket-mode": "^2.0.0",
"@slack/web-api": "^7.0.0",
"@aws-sdk/client-lambda": "^3.1030.0",
"@aws-sdk/client-secrets-manager": "^3.1030.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
"typescript": "~5.7.2"
}
}

View file

@ -0,0 +1,85 @@
import { SocketModeClient } from '@slack/socket-mode';
import { LambdaClient, InvokeCommand } from '@aws-sdk/client-lambda';
import { SecretsManagerClient, GetSecretValueCommand } from '@aws-sdk/client-secrets-manager';
const lambdaClient = new LambdaClient({});
const secretsClient = new SecretsManagerClient({});
const APP_TOKEN_SECRET_ARN = process.env.APP_TOKEN_SECRET_ARN!;
const PROCESSOR_FUNCTION_NAME = process.env.PROCESSOR_FUNCTION_NAME!;
const APP_HOME_FUNCTION_NAME = process.env.APP_HOME_FUNCTION_NAME!;
async function getAppToken(): Promise<string> {
const res = await secretsClient.send(
new GetSecretValueCommand({ SecretId: APP_TOKEN_SECRET_ARN }),
);
return res.SecretString!;
}
async function invokeProcessor(payload: Record<string, unknown>): Promise<void> {
await lambdaClient.send(
new InvokeCommand({
FunctionName: PROCESSOR_FUNCTION_NAME,
InvocationType: 'Event',
Payload: Buffer.from(JSON.stringify(payload)),
}),
);
}
async function main(): Promise<void> {
const appToken = await getAppToken();
const client = new SocketModeClient({ appToken });
client.on('message', async ({ event, ack }) => {
await ack();
if (event.bot_id || event.subtype) return;
if (event.channel_type !== 'im') return;
console.log('Processing DM from', event.user);
await invokeProcessor({
userId: event.user,
channelId: event.channel,
text: event.text ?? '',
ts: event.ts,
threadTs: event.thread_ts,
eventType: 'dm',
});
});
client.on('app_mention', async ({ event, ack }) => {
await ack();
const text = (event.text ?? '').replace(/^<@[A-Z0-9]+>\s*/, '');
console.log('Processing @mention from', event.user, 'in', event.channel);
await invokeProcessor({
userId: event.user,
channelId: event.channel,
text,
ts: event.ts,
threadTs: event.thread_ts,
eventType: 'app_mention',
});
});
client.on('app_home_opened', async ({ event, ack }) => {
await ack();
await lambdaClient.send(
new InvokeCommand({
FunctionName: APP_HOME_FUNCTION_NAME,
InvocationType: 'Event',
Payload: Buffer.from(JSON.stringify({ userId: event.user })),
}),
);
});
await client.start();
console.log('Socket Mode client connected');
}
main().catch((err) => {
console.error('Socket Mode startup failed:', err);
process.exit(1);
});

View file

@ -0,0 +1,18 @@
{
"compilerOptions": {
"target": "ES2022",
"module": "commonjs",
"lib": ["ES2022"],
"outDir": "dist",
"rootDir": "src",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"forceConsistentCasingInFileNames": true,
"resolveJsonModule": true,
"declaration": false,
"sourceMap": false
},
"include": ["src/**/*"],
"exclude": ["node_modules", "dist"]
}

View file

@ -14,5 +14,5 @@
"experimentalDecorators": true,
"skipLibCheck": true
},
"exclude": ["node_modules", "dist", "lambda", "scripts"]
"exclude": ["node_modules", "dist", "lambda", "scripts", "services", "cdk.out"]
}