Merge workorder-ingest into unified procurement repo (#22)

* 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.

* Fix test_local.py import path and remove dead shared/models.py

test_local.py referenced the old lambdas/email_processor path. Updated
to lambdas/wo/email_processor. Removed shared/ directory entirely as
nothing imports from it.

* Escape HTML in both web UI dashboards to prevent XSS

Both Function URLs are public (auth_type=NONE) and render
email-derived content via f-strings. Attacker-crafted emails
could inject scripts. Added html.escape() on all interpolated
values in both PO and WO dashboards.

* Add pagination to WO web UI scan

get_work_orders() only fetched the first 1MB page from DynamoDB.
Loop on LastEvaluatedKey to match the PO web UI pattern.

* Fix esc(None) TypeError and javascript: scheme in PO web UI

Coerce supplier name through `or ""` before escaping to handle
nested None from DynamoDB. Add scheme allowlist on view_order_url
to block javascript:/data: hrefs from LLM-extracted URLs.

* Fix WO render_badge None guard, updated_at slice, and backfill path

Add null guard to WO render_badge matching the PO version. Use
`or ""` before slicing updated_at to handle explicit None values.
Fix backfill_sites.py sys.path to use new lambdas/po/site_extractor.

* Harden WO web UI and fix JS-context XSS in both dashboards

- Use json.dumps for onclick URLs to prevent JS string breakout
- Add .lower() to WO render_badge color lookup matching PO pattern
- Add pagination to get_comments query
- Cap get_work_orders to 500 results matching PO pattern

* Apply ruff formatting to web UI handlers
This commit is contained in:
Adam Moussa 2026-05-12 15:21:06 -04:00 • committed by GitHub
parent c0b276f2bc
commit 5112c1345b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
19 changed files with 925 additions and 126 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:

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

@ -5,8 +5,10 @@ Serves a simple HTML dashboard for viewing purchase orders.
Accessed via Lambda Function URL.
"""
import json
import os
from decimal import Decimal
from html import escape as esc
import boto3
@ -31,7 +33,7 @@ def render_badge(value, color_map):
if not value:
value = "unknown"
color = color_map.get(value.lower(), "#9ca3af")
label = value.replace("_", " ").title()
label = esc(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>'
@ -62,50 +64,51 @@ def fmt_currency(val):
def render_po_detail(po):
po_number = po.get("po_number", "")
po_number = esc(po.get("po_number", ""))
fields = [
("PO Number", po_number),
("Status", render_badge(po.get("po_status", ""), STATUS_COLORS)),
("Email Type", render_badge(po.get("email_type", ""), EMAIL_TYPE_COLORS)),
("Total Amount", fmt_currency(po.get("total_amount"))),
("Currency", po.get("currency")),
("Supplier", (po.get("supplier") or {}).get("name")),
("Site Code", po.get("site_code")),
("State", po.get("state")),
("Trade", po.get("trade")),
("Fiscal Year", po.get("fiscal_year")),
("Coupa Category", po.get("coupa_category")),
("Submitted By", po.get("submitted_by")),
("On Behalf Of", po.get("on_behalf_of")),
("Order Date", po.get("order_date")),
("Revision Date", po.get("revision_date")),
("Payment Terms", po.get("payment_terms")),
("Requisition #", po.get("requisition_number")),
("Department", po.get("department")),
("Data Source", po.get("data_source")),
("Processed At", po.get("processed_at")),
("Currency", esc(po.get("currency", "")) or None),
("Supplier", esc((po.get("supplier") or {}).get("name") or "") or None),
("Site Code", esc(po.get("site_code", "")) or None),
("State", esc(po.get("state", "")) or None),
("Trade", esc(po.get("trade", "")) or None),
("Fiscal Year", esc(po.get("fiscal_year", "")) or None),
("Coupa Category", esc(po.get("coupa_category", "")) or None),
("Submitted By", esc(po.get("submitted_by", "")) or None),
("On Behalf Of", esc(po.get("on_behalf_of", "")) or None),
("Order Date", esc(po.get("order_date", "")) or None),
("Revision Date", esc(po.get("revision_date", "")) or None),
("Payment Terms", esc(po.get("payment_terms", "")) or None),
("Requisition #", esc(po.get("requisition_number", "")) or None),
("Department", esc(po.get("department", "")) or None),
("Data Source", esc(po.get("data_source", "")) or None),
("Processed At", esc(po.get("processed_at", "")) or None),
]
ship_to = po.get("ship_to") or {}
if any(ship_to.values()):
ship_parts = []
if ship_to.get("name"):
ship_parts.append(ship_to["name"])
ship_parts.append(esc(ship_to["name"]))
if ship_to.get("address"):
ship_parts.append(ship_to["address"])
ship_parts.append(esc(ship_to["address"]))
if ship_to.get("location_code"):
ship_parts.append(f"Location: {ship_to['location_code']}")
ship_parts.append(f"Location: {esc(ship_to['location_code'])}")
if ship_to.get("attn"):
ship_parts.append(f"Attn: {ship_to['attn']}")
ship_parts.append(f"Attn: {esc(ship_to['attn'])}")
fields.append(("Ship To", "<br>".join(ship_parts)))
view_url = po.get("view_order_url")
if view_url:
if view_url and view_url.startswith(("https://", "http://")):
escaped_url = esc(view_url, quote=True)
fields.append(
(
"Coupa Link",
f'<a href="{view_url}" target="_blank" style="color:#3b82f6;">View in Coupa</a>',
f'<a href="{escaped_url}" target="_blank" style="color:#3b82f6;">View in Coupa</a>',
)
)
@ -124,17 +127,17 @@ def render_po_detail(po):
if line_items:
rows = ""
for item in line_items:
qty = item.get("quantity", "") or ""
unit = item.get("unit", "") or ""
price = item.get("price", "") or ""
qty = esc(str(item.get("quantity", "") or ""))
unit = esc(str(item.get("unit", "") or ""))
price = esc(str(item.get("price", "") or ""))
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;">
<td style="padding:10px;font-size:14px;">{item.get("description", "")}</td>
<td style="padding:10px;font-size:14px;">{esc(str(item.get("description", "")))}</td>
<td style="padding:10px;font-size:13px;text-align:right;">{qty}</td>
<td style="padding:10px;font-size:13px;">{unit}</td>
<td style="padding:10px;font-size:13px;text-align:right;">{price}</td>
<td style="padding:10px;font-size:14px;text-align:right;">{fmt_currency(item.get("amount"))}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{item.get("need_by", "") or ""}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{esc(str(item.get("need_by", "") or ""))}</td>
</tr>"""
items_html = f"""
@ -184,16 +187,16 @@ def render_po_detail(po):
def render_po_list(purchase_orders):
rows = ""
for po in purchase_orders:
po_number = po.get("po_number", "")
supplier = (po.get("supplier") or {}).get("name", "")
site_code = po.get("site_code", "")
trade = po.get("trade", "")
po_number = esc(po.get("po_number", ""))
supplier = esc((po.get("supplier") or {}).get("name", ""))
site_code = esc(po.get("site_code", ""))
trade = esc(po.get("trade", ""))
status = po.get("po_status", "")
total = fmt_currency(po.get("total_amount"))
processed = (po.get("processed_at") or "")[:16]
processed = esc((po.get("processed_at") or "")[:16])
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location='/po?id={po_number}'">
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location={esc(json.dumps(f"/po?id={po.get("po_number", "")}"), quote=True)}">
<td style="padding:12px;font-weight:500;color:#3b82f6;">{po_number}</td>
<td style="padding:12px;">{supplier}</td>
<td style="padding:12px;font-weight:500;">{site_code}</td>

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,256 @@
"""
Web UI Lambda.
Serves a simple HTML dashboard for viewing work orders and comments.
Accessed via Lambda Function URL.
"""
import json
import os
from html import escape as esc
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(limit=500):
table = dynamodb.Table(WORK_ORDERS_TABLE)
items = []
response = table.scan()
items.extend(response.get("Items", []))
while "LastEvaluatedKey" in response:
response = table.scan(ExclusiveStartKey=response["LastEvaluatedKey"])
items.extend(response.get("Items", []))
items.sort(key=lambda x: x.get("updated_at", ""), reverse=True)
return items[:limit]
def get_comments(work_order_id):
table = dynamodb.Table(COMMENTS_TABLE)
items = []
kwargs = {
"KeyConditionExpression": "work_order_id = :woid",
"ExpressionAttributeValues": {":woid": work_order_id},
}
response = table.query(**kwargs)
items.extend(response.get("Items", []))
while "LastEvaluatedKey" in response:
response = table.query(**kwargs, ExclusiveStartKey=response["LastEvaluatedKey"])
items.extend(response.get("Items", []))
items.sort(key=lambda x: x.get("created_at", ""), reverse=True)
return items
def render_badge(value, color_map):
if not value:
value = "unknown"
color = color_map.get(value.lower(), "#9ca3af")
label = esc(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 = esc(e.get("commenter", ""))
created = esc(e.get("created_at", ""))
text = esc(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 = esc(wo.get("work_order_id", ""))
fields = [
("Description", esc(wo.get("description", "")) or None),
("Status", render_badge(wo.get("wo_status", "unknown"), STATUS_COLORS)),
(
"Record Type",
render_badge(wo.get("record_type", "unknown"), RECORD_TYPE_COLORS),
),
("Site Code", esc(wo.get("site_code", "")) or None),
("Building", esc(wo.get("building", "")) or None),
("Address", esc(wo.get("address", "")) or None),
("Severity", esc(wo.get("severity", "")) or None),
("Priority", esc(wo.get("priority", "")) or None),
("Due Date", esc(wo.get("due_date", "")) or None),
("Date Reported", esc(wo.get("date_reported", "")) or None),
("Scheduled Start", esc(wo.get("scheduled_start", "")) or None),
("Assigned To", esc(wo.get("assigned_to", "")) or None),
("Created", esc(wo.get("created_at", "")) or None),
("Last Updated", esc(wo.get("updated_at", "")) or None),
]
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 = esc(wo.get("work_order_id", ""))
desc = esc(wo.get("description", ""))
site = esc(wo.get("site_code", ""))
status = wo.get("wo_status", "unknown")
record_type = wo.get("record_type", "unknown")
due = esc(wo.get("due_date", ""))
updated = esc((wo.get("updated_at") or "")[:16])
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location={esc(json.dumps(f"/wo?id={wo.get("work_order_id", "")}"), quote=True)}">
<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

View file

@ -8,7 +8,7 @@ Usage:
import sys
import os
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "lambdas", "site_extractor"))
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "lambdas", "po", "site_extractor"))
import boto3
from handler import extract_site_code, parse_address, upsert_site

47
test_local.py Normal file
View file

@ -0,0 +1,47 @@
#!/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
sys.path.insert(0, str(Path(__file__).parent / "lambdas" / "wo" / "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()