Merge workorder-ingest pipeline into unified repo

Move PO lambdas under lambdas/po/, add WO pipeline under lambdas/wo/.
Two independent CloudFormation stacks in one CDK app. Fix WO stack
compliance: ARM64 architecture, 60-day log retention, aarch64 bundling,
RETAIN on Anthropic secret. Remove stale CodePipeline buildspec.
This commit is contained in:
Adam Moussa 2026-05-12 14:12:17 -04:00
parent c0b276f2bc
commit 5a1e774b7c
21 changed files with 959 additions and 91 deletions

View file

@ -10,7 +10,25 @@ updates:
- "minor"
- "patch"
- package-ecosystem: "pip"
directory: "/lambdas/email_processor"
directory: "/lambdas/po/email_processor"
schedule:
interval: "weekly"
groups:
minor-and-patch:
update-types:
- "minor"
- "patch"
- package-ecosystem: "pip"
directory: "/lambdas/wo/email_processor"
schedule:
interval: "weekly"
groups:
minor-and-patch:
update-types:
- "minor"
- "patch"
- package-ecosystem: "pip"
directory: "/lambdas/wo/web_ui"
schedule:
interval: "weekly"
groups:

View file

@ -7,6 +7,6 @@ jobs:
ci:
uses: Sea-Haven-Industries/.github/.github/workflows/ci-python-sam.yaml@main
with:
source-dirs: "lambdas cdk"
source-dirs: "lambdas cdk shared"
run-cdk-synth: true
run-sam-validate: false

2
.gitignore vendored
View file

@ -9,4 +9,4 @@ node_modules/
cdk.out/
.env
*.eml
lambdas/*/package/
lambdas/*/*/package/

159
README.md
View file

@ -1,106 +1,127 @@
# PO Ingest
# Procurement Ingest
Coupa purchase-order email ingestion pipeline. SES receives Amazon PO emails, Claude extracts structured data, and the result lands in the shared `purchase-orders` DynamoDB table.
Unified email ingestion pipelines for Amazon procurement data. Two independent pipelines — purchase orders (Coupa) and work orders (APM/Hexagon EAM) — share a single repo and CDK app but deploy as separate CloudFormation stacks.
## Flow
## Pipelines
1. Coupa sends a PO email to `amazon_po@int.seahaven.com`.
2. SES (using the shared `INBOUND_MAIL` rule set) drops the raw MIME into `s3://po-ingest-emails-{AccountId}/inbound/`.
3. S3 `ObjectCreated` fires the `po-email-processor` Lambda.
4. The Lambda parses the email, sends it to Claude Haiku 4.5 for structured JSON extraction, and writes to DynamoDB.
- `email_type: new_po` — conditional `PutItem` on `purchase-orders` (idempotent on `po_number`).
- `email_type: revision` — unconditional `PutItem` overwriting the existing record with updated data.
- `email_type: cancellation` — `UpdateItem` marking the existing row `Cancelled`.
5. DynamoDB Streams (NEW_IMAGE) on `purchase-orders` feeds two downstream consumers:
- **LedgerFlow** (`seahaven-slack-bot/po-sync`) — daily KB sync.
- **Verified-sites pipeline** (`po-ingest-site-extractor`) — real-time site address extraction (see below).
### Purchase Orders (`po-ingest` stack)
The extraction prompt includes domain-specific rules for site code identification (with a skip list for false positives like RME, BBM, JLL), trade classification across 23 categories (Plumbing PM/Reactive, Electrical, HVAC, Dock Doors, etc.), fiscal year derivation, and ship-to address parsing with zip code zero-padding.
Coupa PO emails are received at `amazon_po@int.seahaven.com`, parsed by Claude Haiku 4.5, and written to the `purchase-orders` DynamoDB table.
A separate `po-web-ui` Lambda (Function URL, unauthenticated) renders a simple HTML dashboard scanning the table.
**Flow:**
1. Coupa sends a PO email (new, revision, or cancellation).
2. SES (`INBOUND_MAIL` rule set) drops the raw MIME into `s3://po-ingest-emails-{AccountId}/inbound/`.
3. S3 `ObjectCreated` triggers the `po-email-processor` Lambda.
4. Claude extracts structured JSON (PO number, status, supplier, site code, trade classification, line items, fiscal year).
5. Conditional write to DynamoDB:
- `new_po` — idempotent insert (no-op if PO exists)
- `revision` — unconditional overwrite
- `cancellation` — marks existing row `Cancelled`
6. DynamoDB Streams feeds downstream consumers:
- **LedgerFlow** (`seahaven-slack-bot/po-sync`) — daily KB sync
- **Site extractor** (`po-ingest-site-extractor`) — real-time site address extraction into `verified-sites` table
### PO record schema
**Lambdas** (`lambdas/po/`):
| Function | Trigger | Purpose |
|---|---|---|
| `po-email-processor` | S3 ObjectCreated | Claude extraction + DynamoDB write |
| `po-ingest-site-extractor` | DynamoDB Streams | Site code/address extraction -> `verified-sites` |
| `po-web-ui` | Function URL | HTML dashboard |
Each record in `purchase-orders` includes:
- **Core**: `po_number` (PK), `email_type` (new_po/revision/cancellation), `po_status`, `source_system`
- **People/dates**: `submitted_by`, `on_behalf_of`, `order_date`, `revision_date`, `payment_terms`, `requisition_number`, `department`
- **Site**: `site_code`, `state` (top-level), `ship_to` (structured), `ship_to_raw` (original text)
- **Classification**: `trade`, `fiscal_year`, `coupa_category`
- **Financials**: `total_amount`, `currency`, `line_items[]` (with `description`, `amount`, `quantity`, `unit`, `price`, `need_by`)
- **Metadata**: `data_source` ("email" or "email+payee_scrape"), `email_subject`, `processed_at`, `raw_s3_key`
**Tables:**
- `purchase-orders` (PK: `po_number`, Streams: NEW_IMAGE) — shared with payments-dashboard and seahaven-slack-bot
- `verified-sites` (PK: `siteCode`, GSI: `by-state`) — ~1,100 unique Amazon facility sites
- `pending-site-review` (PK: `po_number`) — unresolvable POs for manual Payee Central verification
### Verified-sites pipeline
### Work Orders (`WorkorderIngestStack` stack)
The `po-ingest-site-extractor` Lambda is triggered by the DynamoDB Stream on every PO INSERT/MODIFY. It:
Amazon APM work order emails (from Hexagon EAM / HxGN SmartCloud) are received at `apm@int.seahaven.com`, parsed by Claude Haiku 4.5, and written to the `WorkOrders` DynamoDB table.
1. Extracts an Amazon facility site code from `ship_to.name` using a regex cascade (parentheses, `LLC - CODE`, `Station CODE`, `DS - CODE`) with a fallback to the first `line_items` description.
2. Parses `ship_to.address` into structured fields (street, city, state, zip).
3. Upserts to the `verified-sites` DynamoDB table — atomically increments `poCount` and appends the PO number to `sourcePOs`.
**Flow:**
1. Hexagon EAM sends email notifications (new assignments, comments, updates, cancellations) to `amazon@seahavenind.com`.
2. Gmail filter forwards APM emails to `apm@int.seahaven.com` (SES).
3. SES drops the raw MIME into `s3://workorder-ingest-emails-{AccountId}/inbound/`.
4. S3 triggers the `workorder-email-processor` Lambda.
5. Claude extracts structured JSON (work order ID, site code, severity, priority, dates, assigned technician).
6. Work order upserted to `WorkOrders`, event/comment appended to `WorkOrderComments`.
POs with no extractable site code fall through to an address reverse-lookup against the verified-sites cache (normalized street + zip). If still unresolved, the PO is written to the `pending-site-review` table for manual verification against Payee Central.
**Lambdas** (`lambdas/wo/`):
| Function | Trigger | Purpose |
|---|---|---|
| `workorder-email-processor` | S3 ObjectCreated | Claude extraction + DynamoDB write |
| `workorder-web-ui` | Function URL | HTML dashboard |
Backfill stats (initial run): 14,825 POs scanned → 9,900 with extractable site codes → 1,100 unique sites.
**Tables:**
- `WorkOrders` (PK: `work_order_id`, GSIs: `site-code-index`, `status-index`)
- `WorkOrderComments` (PK: `work_order_id`, SK: `comment_id`)
## Architecture
- **IaC:** AWS CDK (Python), stack name `po-ingest`, region `us-east-1`.
- **Lambdas** (all Python 3.12, arm64, 60-day log retention):
- `po-email-processor` — S3-triggered, parses PO emails via Claude Haiku.
- `po-web-ui` — Function URL, HTML dashboard.
- `po-ingest-site-extractor` — DynamoDB Streams-triggered, extracts site addresses.
- **Storage:**
- S3 `po-ingest-emails-{AccountId}` — 90-day lifecycle expiry.
- DynamoDB `purchase-orders` — owned by this stack, Streams enabled (NEW_AND_OLD_IMAGES).
- DynamoDB `verified-sites` — PK `siteCode`, GSI `by-state` on `state`.
- DynamoDB `pending-site-review` — PK `po_number`. POs with no extractable site code and no address match, awaiting manual Payee Central verification.
- **Secrets:** Anthropic API key in Secrets Manager at `po-ingest/anthropic-api-key`.
- **SES:** adds the `PoEmailRule` to the existing `INBOUND_MAIL` receipt rule set (shared with `workorder-ingest`).
- **CI/CD:** CodePipeline V2 (`po-ingest-pipeline`) → CodeBuild (`po-ingest-build`). Pushes to `main` auto-deploy via `buildspec.yml`.
**IaC:** AWS CDK (Python), two stacks in one app, region `us-east-1`.
All Lambdas: Python 3.12, ARM64, 60-day log retention.
**Secrets:**
- `po-ingest/anthropic-api-key` — Anthropic API key for PO parsing
- `workorder-ingest/anthropic-api-key` — Anthropic API key for WO parsing
**SES:** Both stacks add rules to the shared `INBOUND_MAIL` receipt rule set on `int.seahaven.com`.
## CI/CD
Merges to `main` trigger the `po-ingest-pipeline` (CodePipeline V2) which runs CodeBuild to `cdk deploy`. The pipeline uses the existing CodeStar connection to the Sea-Haven-Industries GitHub org.
GitHub Actions with reusable workflows from `Sea-Haven-Industries/.github`:
- **CI** (PR to `main`): linting + `cdk synth` via `ci-python-sam.yaml@main`
- **CD** (push to `main`): `cdk deploy --all` via `cd-cdk.yaml@main` (OIDC auth)
**Branch protection:** `main` requires a PR (no direct push), no deletion, no force push.
Branch protection on `main` — all changes through PR.
## Setup
1. Bootstrap CDK in the account if you haven't already: `cdk bootstrap aws://{AccountId}/us-east-1`.
2. Store the Anthropic API key:
1. Bootstrap CDK: `cdk bootstrap aws://{AccountId}/us-east-1`
2. Store Anthropic API keys:
```bash
aws secretsmanager create-secret \
--name po-ingest/anthropic-api-key \
--secret-string "sk-ant-..."
aws secretsmanager create-secret --name po-ingest/anthropic-api-key --secret-string "sk-ant-..."
aws secretsmanager create-secret --name workorder-ingest/anthropic-api-key --secret-string "sk-ant-..."
```
3. Install Lambda dependencies into the deployable package directory (gitignored):
```bash
pip install -r lambdas/email_processor/requirements.txt -t lambdas/email_processor/package/
```
4. Deploy:
3. Deploy both stacks:
```bash
cd cdk
pip install -r requirements.txt
cdk deploy
cdk deploy --all
```
5. The `WebUIUrl` CloudFormation output is the dashboard URL.
4. CloudFormation outputs include `WebUIUrl` for each stack's dashboard.
## Reprocessing
To re-run the processor against every email still sitting in `inbound/` (useful after a parser change):
## Scripts
**Reprocess PO emails** (re-run parser against all emails still in S3):
```bash
python scripts/reprocess.py # dry-run — lists keys
python scripts/reprocess.py --execute # invokes po-email-processor for each
python scripts/reprocess.py # dry-run
python scripts/reprocess.py --execute # invoke po-email-processor for each
```
Inserts are conditional on `po_number`, so re-processing existing POs is a no-op.
## Backfilling verified sites
The stream Lambda handles all future POs automatically. To backfill from historical PO data (one-time):
**Backfill verified sites** (one-time scan of historical POs):
```bash
python scripts/backfill_sites.py
```
Uses the same extraction logic as the Lambda. Idempotent — safe to re-run.
## Directory Structure
```
cdk/
app.py # Two stacks: po-ingest + WorkorderIngestStack
po_stack.py # Purchase order pipeline resources
wo_stack.py # Work order pipeline resources
lambdas/
po/ # PO pipeline Lambdas
email_processor/
site_extractor/
web_ui/
wo/ # WO pipeline Lambdas
email_processor/
web_ui/
shared/
models.py # Work order dataclasses/enums
scripts/
reprocess.py
backfill_sites.py
```

View file

@ -1,15 +0,0 @@
version: 0.2
phases:
install:
runtime-versions:
python: 3.12
nodejs: 22
commands:
- npm install -g aws-cdk
- pip install -r cdk/requirements.txt
- pip install -r lambdas/email_processor/requirements.txt -t lambdas/email_processor/package/
- cp lambdas/email_processor/handler.py lambdas/email_processor/package/
build:
commands:
- cd cdk && cdk deploy --require-approval never

View file

@ -1,12 +1,22 @@
#!/usr/bin/env python3
import aws_cdk as cdk
from stack import PoIngestStack
from po_stack import PoIngestStack
from wo_stack import WorkorderIngestStack
app = cdk.App()
PoIngestStack(
app,
"po-ingest",
stack_name="po-ingest",
env=cdk.Environment(region="us-east-1"),
)
WorkorderIngestStack(
app,
"workorder-ingest",
stack_name="WorkorderIngestStack",
env=cdk.Environment(region="us-east-1"),
)
app.synth()

View file

@ -67,7 +67,7 @@ class PoIngestStack(Stack):
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset(
"../lambdas/email_processor",
"../lambdas/po/email_processor",
bundling=cdk.BundlingOptions(
image=lambda_.Runtime.PYTHON_3_12.bundling_image,
command=[
@ -125,7 +125,7 @@ class PoIngestStack(Stack):
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/web_ui"),
code=lambda_.Code.from_asset("../lambdas/po/web_ui"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
@ -174,7 +174,7 @@ class PoIngestStack(Stack):
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/site_extractor"),
code=lambda_.Code.from_asset("../lambdas/po/site_extractor"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,

185
cdk/wo_stack.py Normal file
View file

@ -0,0 +1,185 @@
"""CDK stack for the work order email ingestion pipeline."""
import aws_cdk as cdk
from aws_cdk import (
Duration,
RemovalPolicy,
Stack,
aws_dynamodb as dynamodb,
aws_lambda as lambda_,
aws_logs as logs,
aws_s3 as s3,
aws_s3_notifications as s3n,
aws_ses as ses,
aws_ses_actions as ses_actions,
aws_secretsmanager as secretsmanager,
)
from constructs import Construct
class WorkorderIngestStack(Stack):
def __init__(self, scope: Construct, construct_id: str, **kwargs):
super().__init__(scope, construct_id, **kwargs)
# --- S3 bucket for raw emails ---
email_bucket = s3.Bucket(
self,
"EmailBucket",
bucket_name=f"workorder-ingest-emails-{self.account}",
removal_policy=RemovalPolicy.RETAIN,
lifecycle_rules=[
s3.LifecycleRule(expiration=Duration.days(90)),
],
)
# --- DynamoDB tables ---
work_orders_table = dynamodb.Table(
self,
"WorkOrdersTable",
table_name="WorkOrders",
partition_key=dynamodb.Attribute(
name="work_order_id",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
work_orders_table.add_global_secondary_index(
index_name="site-code-index",
partition_key=dynamodb.Attribute(
name="site_code",
type=dynamodb.AttributeType.STRING,
),
sort_key=dynamodb.Attribute(
name="updated_at",
type=dynamodb.AttributeType.STRING,
),
)
work_orders_table.add_global_secondary_index(
index_name="status-index",
partition_key=dynamodb.Attribute(
name="wo_status",
type=dynamodb.AttributeType.STRING,
),
sort_key=dynamodb.Attribute(
name="updated_at",
type=dynamodb.AttributeType.STRING,
),
)
comments_table = dynamodb.Table(
self,
"CommentsTable",
table_name="WorkOrderComments",
partition_key=dynamodb.Attribute(
name="work_order_id",
type=dynamodb.AttributeType.STRING,
),
sort_key=dynamodb.Attribute(
name="comment_id",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
# --- Secrets Manager for Anthropic API key ---
anthropic_secret = secretsmanager.Secret(
self,
"AnthropicApiKey",
secret_name="workorder-ingest/anthropic-api-key",
description="Anthropic API key for work order email parsing",
removal_policy=RemovalPolicy.RETAIN,
)
# --- Lambda function ---
email_processor = lambda_.Function(
self,
"EmailProcessor",
function_name="workorder-email-processor",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset(
"../lambdas/wo/email_processor",
bundling=cdk.BundlingOptions(
image=lambda_.Runtime.PYTHON_3_12.bundling_image,
command=[
"bash",
"-c",
"pip install --platform manylinux2014_aarch64 --only-binary=:all: "
"-r requirements.txt -t /asset-output && "
"cp -r . /asset-output/",
],
),
),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={
"WORK_ORDERS_TABLE": work_orders_table.table_name,
"COMMENTS_TABLE": comments_table.table_name,
"ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn,
},
)
# Grant permissions
email_bucket.grant_read(email_processor)
work_orders_table.grant_read_write_data(email_processor)
comments_table.grant_read_write_data(email_processor)
anthropic_secret.grant_read(email_processor)
# S3 event notification -> Lambda
email_bucket.add_event_notification(
s3.EventType.OBJECT_CREATED,
s3n.LambdaDestination(email_processor),
s3.NotificationKeyFilter(prefix="inbound/"),
)
# --- SES Receipt Rule ---
rule_set = ses.ReceiptRuleSet.from_receipt_rule_set_name(
self,
"ExistingRuleSet",
"INBOUND_MAIL",
)
rule_set.add_rule(
"WorkorderEmailRule",
recipients=["apm@int.seahaven.com"],
actions=[
ses_actions.S3(
bucket=email_bucket,
object_key_prefix="inbound/",
),
],
)
# --- Web UI Lambda ---
web_ui = lambda_.Function(
self,
"WebUI",
function_name="workorder-web-ui",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/wo/web_ui"),
timeout=Duration.seconds(15),
memory_size=128,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={
"WORK_ORDERS_TABLE": work_orders_table.table_name,
"COMMENTS_TABLE": comments_table.table_name,
},
)
work_orders_table.grant_read_data(web_ui)
comments_table.grant_read_data(web_ui)
# Function URL for direct access
web_url = web_ui.add_function_url(
auth_type=lambda_.FunctionUrlAuthType.NONE,
)
cdk.CfnOutput(
self, "WebUIUrl", value=web_url.url, description="Work Order Dashboard URL"
)

View file

View file

@ -0,0 +1,271 @@
"""
Email processor Lambda.
Triggered by S3 events when SES delivers an email.
Parses the raw email, sends it to Claude for structured extraction,
then writes the result to DynamoDB.
"""
import email
import json
import logging
import os
import re
from datetime import datetime
from email import policy
import anthropic
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
s3 = boto3.client("s3")
dynamodb = boto3.resource("dynamodb")
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN")
EXTRACTION_PROMPT = """\
You are an email parser for a facilities maintenance work order system.
The emails come from Amazon's APM system (via Hexagon EAM / HxGN SmartCloud).
Analyze the following email and extract structured data. Return ONLY valid JSON with these fields:
{
"email_type": "new_work_order" | "update" | "comment" | "cancellation",
"work_order_id": "string or null",
"description": "work order description or null",
"status": "new" | "assigned" | "in_progress" | "on_hold" | "completed" | "cancelled" | "unknown",
"site_code": "building/site code like WIL1, ZDL8, etc. or null",
"building": "full building identifier or null",
"address": "physical address or null",
"severity": "severity level or null",
"priority": "priority level or null",
"date_reported": "ISO 8601 date or null",
"scheduled_start": "ISO 8601 date or null",
"due_date": "ISO 8601 date or null",
"assigned_to": "person/team assigned or null",
"commenter": "person who left a comment or null",
"comment_text": "the comment text or null",
"comment_time": "ISO 8601 datetime of the comment or null"
}
Rules:
- "email_type" detection:
- "new_work_order": email announces a new WO assignment
- "comment": email contains a new comment on an existing WO
- "cancellation": email announces a WO has been cancelled
- "update": any other update to an existing WO (status change, reassignment, etc.)
- Extract the site_code from the building field (e.g., "WIL1" from "building WIL1")
- Dates should be converted to ISO 8601 format
- If a field is not present in the email, set it to null
- Do NOT invent or infer data that is not explicitly in the email
"""
def get_anthropic_client() -> anthropic.Anthropic:
"""Create Anthropic client, fetching API key from Secrets Manager if configured."""
if ANTHROPIC_API_KEY_SECRET_ARN:
secrets = boto3.client("secretsmanager")
secret = secrets.get_secret_value(SecretId=ANTHROPIC_API_KEY_SECRET_ARN)
api_key = secret["SecretString"]
return anthropic.Anthropic(api_key=api_key)
# Fall back to ANTHROPIC_API_KEY env var (for local testing)
return anthropic.Anthropic()
def parse_raw_email(raw_bytes: bytes) -> dict:
"""Parse a raw email into subject, sender, body text."""
msg = email.message_from_bytes(raw_bytes, policy=policy.default)
subject = msg.get("Subject", "")
sender = msg.get("From", "")
to = msg.get("To", "")
cc = msg.get("Cc", "")
date = msg.get("Date", "")
body = ""
if msg.is_multipart():
for part in msg.walk():
content_type = part.get_content_type()
if content_type == "text/plain":
body = part.get_content()
break
elif content_type == "text/html" and not body:
body = part.get_content()
else:
body = msg.get_content()
return {
"subject": subject,
"sender": sender,
"to": to,
"cc": cc,
"date": date,
"body": body,
}
def extract_with_claude(email_data: dict) -> dict:
"""Send parsed email to Claude for structured extraction."""
client = get_anthropic_client()
email_text = (
f"Subject: {email_data['subject']}\n"
f"From: {email_data['sender']}\n"
f"To: {email_data['to']}\n"
f"CC: {email_data['cc']}\n"
f"Date: {email_data['date']}\n"
f"\n---\n\n"
f"{email_data['body']}"
)
response = client.messages.create(
model="claude-haiku-4-5-20251001",
max_tokens=1024,
messages=[
{
"role": "user",
"content": f"{EXTRACTION_PROMPT}\n\nEMAIL:\n{email_text}",
}
],
)
response_text = response.content[0].text
# Extract JSON from response (handle markdown code blocks)
json_match = re.search(r"```(?:json)?\s*(.*?)```", response_text, re.DOTALL)
if json_match:
response_text = json_match.group(1)
return json.loads(response_text.strip())
def save_work_order(parsed: dict, s3_key: str):
"""Create or update a work order in DynamoDB."""
table = dynamodb.Table(WORK_ORDERS_TABLE)
work_order_id = parsed["work_order_id"]
now = datetime.utcnow().isoformat()
# Build update expression dynamically from non-null fields
field_map = {
"description": "description",
"status": "wo_status", # 'status' is a DynamoDB reserved word
"site_code": "site_code",
"building": "building",
"address": "address",
"severity": "severity",
"priority": "priority",
"date_reported": "date_reported",
"scheduled_start": "scheduled_start",
"due_date": "due_date",
"assigned_to": "assigned_to",
}
update_parts = ["#updated_at = :updated_at", "#source_key = :source_key"]
attr_names = {
"#updated_at": "updated_at",
"#source_key": "source_email_s3_key",
}
attr_values = {
":updated_at": now,
":source_key": s3_key,
}
for src_field, dynamo_field in field_map.items():
value = parsed.get(src_field)
if value is not None:
placeholder = f":{dynamo_field}"
name_placeholder = f"#{dynamo_field}"
update_parts.append(f"{name_placeholder} = {placeholder}")
attr_names[name_placeholder] = dynamo_field
attr_values[placeholder] = value
# For new items, set created_at
update_parts.append("#created_at = if_not_exists(#created_at, :created_at)")
attr_names["#created_at"] = "created_at"
attr_values[":created_at"] = now
# Customer is always AMAZON for now
update_parts.append("#customer = :customer")
attr_names["#customer"] = "customer"
attr_values[":customer"] = "AMAZON"
# Track the record type (new_work_order, update, comment)
email_type = parsed.get("email_type")
if email_type:
update_parts.append("#record_type = :record_type")
attr_names["#record_type"] = "record_type"
attr_values[":record_type"] = email_type
table.update_item(
Key={"work_order_id": work_order_id},
UpdateExpression="SET " + ", ".join(update_parts),
ExpressionAttributeNames=attr_names,
ExpressionAttributeValues=attr_values,
)
logger.info(f"Saved work order {work_order_id}")
def save_event(parsed: dict, s3_key: str):
"""Save an event to the events table. Every email creates an event entry."""
table = dynamodb.Table(COMMENTS_TABLE)
work_order_id = parsed["work_order_id"]
email_type = parsed.get("email_type", "unknown")
event_time = parsed.get("comment_time") or datetime.utcnow().isoformat()
event_id = f"{work_order_id}#{event_time}"
item = {
"work_order_id": work_order_id,
"comment_id": event_id, # keeping key name for table compatibility
"record_type": email_type,
"commenter": parsed.get("commenter") or "",
"text": parsed.get("comment_text") or "",
"created_at": event_time,
"source_email_s3_key": s3_key,
"ingested_at": datetime.utcnow().isoformat(),
}
table.put_item(Item=item)
logger.info(f"Saved event {event_id} (type={email_type})")
def handler(event, context):
"""Lambda entry point. Triggered by S3 ObjectCreated events."""
for record in event.get("Records", []):
bucket = record["s3"]["bucket"]["name"]
key = record["s3"]["object"]["key"]
logger.info(f"Processing email: s3://{bucket}/{key}")
# Fetch raw email from S3
response = s3.get_object(Bucket=bucket, Key=key)
raw_email = response["Body"].read()
# Parse the raw email
email_data = parse_raw_email(raw_email)
logger.info(f"Subject: {email_data['subject']}")
# Extract structured data with Claude
parsed = extract_with_claude(email_data)
logger.info(
f"Parsed: type={parsed.get('email_type')}, wo={parsed.get('work_order_id')}"
)
if not parsed.get("work_order_id"):
logger.warning(f"No work order ID found in email, skipping: {key}")
continue
s3_key = f"s3://{bucket}/{key}"
# Always upsert the work order with any new info
save_work_order(parsed, s3_key)
# Save every email as an event for history tracking
save_event(parsed, s3_key)
return {"statusCode": 200, "body": "OK"}

View file

@ -0,0 +1,2 @@
anthropic>=0.42.0
boto3>=1.35.0

View file

View file

@ -0,0 +1,244 @@
"""
Web UI Lambda.
Serves a simple HTML dashboard for viewing work orders and comments.
Accessed via Lambda Function URL.
"""
import os
import boto3
dynamodb = boto3.resource("dynamodb")
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
def get_work_orders():
table = dynamodb.Table(WORK_ORDERS_TABLE)
response = table.scan()
items = response.get("Items", [])
# Sort by updated_at descending
items.sort(key=lambda x: x.get("updated_at", ""), reverse=True)
return items
def get_comments(work_order_id):
table = dynamodb.Table(COMMENTS_TABLE)
response = table.query(
KeyConditionExpression="work_order_id = :woid",
ExpressionAttributeValues={":woid": work_order_id},
)
items = response.get("Items", [])
items.sort(key=lambda x: x.get("created_at", ""), reverse=True)
return items
def render_badge(value, color_map):
color = color_map.get(value, "#9ca3af")
label = value.replace("_", " ").title()
return f'<span style="background:{color};color:#fff;padding:2px 10px;border-radius:12px;font-size:12px;font-weight:500;">{label}</span>'
STATUS_COLORS = {
"new": "#3b82f6",
"assigned": "#8b5cf6",
"in_progress": "#f59e0b",
"on_hold": "#6b7280",
"completed": "#10b981",
"cancelled": "#ef4444",
"unknown": "#9ca3af",
}
RECORD_TYPE_COLORS = {
"new_work_order": "#3b82f6",
"comment": "#8b5cf6",
"update": "#f59e0b",
"cancellation": "#ef4444",
"unknown": "#9ca3af",
}
def render_work_order_detail(wo, events):
events_html = ""
if events:
for e in events:
record_type = e.get("record_type", "unknown")
commenter = e.get("commenter", "")
created = e.get("created_at", "")
text = e.get("text", "")
badge = render_badge(record_type, RECORD_TYPE_COLORS)
commenter_str = (
f"<strong>{commenter}</strong> &mdash; " if commenter else ""
)
border_colors = {
"new_work_order": "#3b82f6",
"comment": "#8b5cf6",
"update": "#f59e0b",
"cancellation": "#ef4444",
}
border = border_colors.get(record_type, "#94a3b8")
events_html += f"""
<div style="border-left:3px solid {border};padding:8px 12px;margin-bottom:10px;background:#f8fafc;border-radius:0 6px 6px 0;">
<div style="font-size:12px;color:#64748b;margin-bottom:4px;">
{badge} {commenter_str}{created}
</div>
{"<div style='font-size:14px;color:#1e293b;margin-top:6px;'>" + text + "</div>" if text else ""}
</div>"""
else:
events_html = '<p style="color:#94a3b8;font-style:italic;">No events yet.</p>'
wo_id = wo.get("work_order_id", "")
fields = [
("Description", wo.get("description")),
("Status", render_badge(wo.get("wo_status", "unknown"), STATUS_COLORS)),
(
"Record Type",
render_badge(wo.get("record_type", "unknown"), RECORD_TYPE_COLORS),
),
("Site Code", wo.get("site_code")),
("Building", wo.get("building")),
("Address", wo.get("address")),
("Severity", wo.get("severity")),
("Priority", wo.get("priority")),
("Due Date", wo.get("due_date")),
("Date Reported", wo.get("date_reported")),
("Scheduled Start", wo.get("scheduled_start")),
("Assigned To", wo.get("assigned_to")),
("Created", wo.get("created_at")),
("Last Updated", wo.get("updated_at")),
]
details_html = ""
for label, value in fields:
if value:
details_html += f"""
<div style="display:flex;padding:8px 0;border-bottom:1px solid #f1f5f9;">
<div style="width:140px;font-size:13px;color:#64748b;font-weight:500;">{label}</div>
<div style="flex:1;font-size:14px;color:#1e293b;">{value}</div>
</div>"""
return f"""<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>WO {wo_id} - Sea Haven</title>
<style>
* {{ margin: 0; padding: 0; box-sizing: border-box; }}
body {{ font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', sans-serif; background: #f1f5f9; color: #1e293b; }}
</style>
</head>
<body>
<div style="max-width:800px;margin:0 auto;padding:20px;">
<div style="margin-bottom:20px;">
<a href="/" style="color:#3b82f6;text-decoration:none;font-size:14px;">&larr; All Work Orders</a>
</div>
<div style="background:#fff;border-radius:10px;padding:24px;box-shadow:0 1px 3px rgba(0,0,0,0.08);margin-bottom:20px;">
<h1 style="font-size:20px;margin-bottom:16px;">Work Order {wo_id}</h1>
{details_html}
</div>
<div style="background:#fff;border-radius:10px;padding:24px;box-shadow:0 1px 3px rgba(0,0,0,0.08);">
<h2 style="font-size:16px;margin-bottom:16px;">Events ({len(events)})</h2>
{events_html}
</div>
</div>
</body>
</html>"""
def render_work_orders_list(work_orders):
rows = ""
for wo in work_orders:
wo_id = wo.get("work_order_id", "")
desc = wo.get("description", "")
site = wo.get("site_code", "")
status = wo.get("wo_status", "unknown")
record_type = wo.get("record_type", "unknown")
due = wo.get("due_date", "")
updated = wo.get("updated_at", "")[:16]
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location='/wo?id={wo_id}'">
<td style="padding:12px;font-weight:500;color:#3b82f6;">{wo_id}</td>
<td style="padding:12px;max-width:250px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap;">{desc}</td>
<td style="padding:12px;">{site}</td>
<td style="padding:12px;">{render_badge(record_type, RECORD_TYPE_COLORS)}</td>
<td style="padding:12px;">{render_badge(status, STATUS_COLORS)}</td>
<td style="padding:12px;">{due}</td>
<td style="padding:12px;color:#64748b;font-size:13px;">{updated}</td>
</tr>"""
return f"""<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Work Orders - Sea Haven</title>
<style>
* {{ margin: 0; padding: 0; box-sizing: border-box; }}
body {{ font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', sans-serif; background: #f1f5f9; color: #1e293b; }}
table {{ width: 100%; border-collapse: collapse; }}
tr:hover {{ background: #f8fafc; }}
th {{ text-align: left; padding: 12px; font-size: 12px; text-transform: uppercase; color: #64748b; border-bottom: 2px solid #e2e8f0; }}
</style>
</head>
<body>
<div style="max-width:1100px;margin:0 auto;padding:20px;">
<div style="display:flex;justify-content:space-between;align-items:center;margin-bottom:20px;">
<h1 style="font-size:22px;">Work Orders</h1>
<span style="color:#64748b;font-size:14px;">{len(work_orders)} total</span>
</div>
<div style="background:#fff;border-radius:10px;box-shadow:0 1px 3px rgba(0,0,0,0.08);overflow:hidden;">
<table>
<thead>
<tr>
<th>WO #</th>
<th>Description</th>
<th>Site</th>
<th>Last Action</th>
<th>Status</th>
<th>Due Date</th>
<th>Updated</th>
</tr>
</thead>
<tbody>
{rows if rows else '<tr><td colspan="7" style="padding:40px;text-align:center;color:#94a3b8;">No work orders yet.</td></tr>'}
</tbody>
</table>
</div>
</div>
</body>
</html>"""
def handler(event, context):
path = event.get("rawPath", "/")
qs = event.get("queryStringParameters") or {}
if path == "/wo" and "id" in qs:
wo_id = qs["id"]
# Get work order
table = dynamodb.Table(WORK_ORDERS_TABLE)
result = table.get_item(Key={"work_order_id": wo_id})
wo = result.get("Item")
if not wo:
return {
"statusCode": 404,
"headers": {"Content-Type": "text/html"},
"body": "<h1>Work order not found</h1>",
}
events = get_comments(wo_id)
html = render_work_order_detail(wo, events)
else:
work_orders = get_work_orders()
html = render_work_orders_list(work_orders)
return {
"statusCode": 200,
"headers": {"Content-Type": "text/html"},
"body": html,
}

View file

@ -0,0 +1 @@
boto3>=1.35.0

0
shared/__init__.py Normal file
View file

83
shared/models.py Normal file
View file

@ -0,0 +1,83 @@
"""Data models for work order ingestion."""
from dataclasses import dataclass, field, asdict
from datetime import datetime
from enum import Enum
from typing import Optional
class EmailType(str, Enum):
NEW_WORK_ORDER = "new_work_order"
UPDATE = "update"
COMMENT = "comment"
UNKNOWN = "unknown"
class WorkOrderStatus(str, Enum):
NEW = "new"
ASSIGNED = "assigned"
IN_PROGRESS = "in_progress"
ON_HOLD = "on_hold"
COMPLETED = "completed"
CANCELLED = "cancelled"
UNKNOWN = "unknown"
@dataclass
class Comment:
work_order_id: str
comment_id: str # generated: {work_order_id}#{timestamp}
commenter: str
text: str
created_at: str # ISO 8601
source_email_s3_key: str
ingested_at: str = field(default_factory=lambda: datetime.utcnow().isoformat())
def to_dynamo_item(self) -> dict:
return {k: v for k, v in asdict(self).items() if v is not None}
@dataclass
class WorkOrder:
work_order_id: str
description: Optional[str] = None
status: str = WorkOrderStatus.UNKNOWN.value
customer: str = "AMAZON"
site_code: Optional[str] = None
building: Optional[str] = None
address: Optional[str] = None
severity: Optional[str] = None
priority: Optional[str] = None
date_reported: Optional[str] = None
scheduled_start: Optional[str] = None
due_date: Optional[str] = None
assigned_to: Optional[str] = None
source_email_s3_key: Optional[str] = None
created_at: str = field(default_factory=lambda: datetime.utcnow().isoformat())
updated_at: str = field(default_factory=lambda: datetime.utcnow().isoformat())
def to_dynamo_item(self) -> dict:
return {k: v for k, v in asdict(self).items() if v is not None}
@dataclass
class ParsedEmail:
"""Result of AI parsing an inbound email."""
email_type: str # EmailType value
work_order_id: Optional[str] = None
description: Optional[str] = None
status: Optional[str] = None
site_code: Optional[str] = None
building: Optional[str] = None
address: Optional[str] = None
severity: Optional[str] = None
priority: Optional[str] = None
date_reported: Optional[str] = None
scheduled_start: Optional[str] = None
due_date: Optional[str] = None
assigned_to: Optional[str] = None
commenter: Optional[str] = None
comment_text: Optional[str] = None
comment_time: Optional[str] = None
raw_subject: Optional[str] = None

48
test_local.py Normal file
View file

@ -0,0 +1,48 @@
#!/usr/bin/env python3
"""
Local test script - parses sample emails through Claude without AWS.
Usage:
export ANTHROPIC_API_KEY=sk-ant-...
python test_local.py
"""
import json
import sys
from pathlib import Path
# Add project root to path so we can import the handler's parsing logic
sys.path.insert(0, str(Path(__file__).parent / "lambdas" / "email_processor"))
from handler import parse_raw_email, extract_with_claude
def main():
samples_dir = Path(__file__).parent / "samples"
eml_files = list(samples_dir.glob("*.eml"))
if not eml_files:
print("No .eml files found in samples/")
return
for eml_path in eml_files:
print(f"\n{'='*80}")
print(f"FILE: {eml_path.name}")
print(f"{'='*80}")
raw = eml_path.read_bytes()
email_data = parse_raw_email(raw)
print(f"Subject: {email_data['subject']}")
print(f"From: {email_data['sender']}")
print(f"Date: {email_data['date']}")
print(f"Body preview: {email_data['body'][:200]}...")
print()
print("Sending to Claude for extraction...")
parsed = extract_with_claude(email_data)
print(json.dumps(parsed, indent=2))
if __name__ == "__main__":
main()