Compare commits

..

No commits in common. "cb2af81615be65ec02de7f039938bf47ad3590f3" and "c0b276f2bcded25ab1874fd43f646c08eb187cc0" have entirely different histories.

19 changed files with 126 additions and 925 deletions

View file

@ -10,25 +10,7 @@ updates:
- "minor"
- "patch"
- package-ecosystem: "pip"
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"
directory: "/lambdas/email_processor"
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,127 +1,106 @@
# Procurement Ingest
# PO Ingest
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.
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.
## Pipelines
## Flow
### Purchase Orders (`po-ingest` stack)
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).
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.
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.
**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
A separate `po-web-ui` Lambda (Function URL, unauthenticated) renders a simple HTML dashboard scanning the table.
**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 |
### PO record schema
**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
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`
### Work Orders (`WorkorderIngestStack` stack)
### Verified-sites pipeline
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.
The `po-ingest-site-extractor` Lambda is triggered by the DynamoDB Stream on every PO INSERT/MODIFY. It:
**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`.
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`.
**Lambdas** (`lambdas/wo/`):
| Function | Trigger | Purpose |
|---|---|---|
| `workorder-email-processor` | S3 ObjectCreated | Claude extraction + DynamoDB write |
| `workorder-web-ui` | Function URL | HTML dashboard |
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.
**Tables:**
- `WorkOrders` (PK: `work_order_id`, GSIs: `site-code-index`, `status-index`)
- `WorkOrderComments` (PK: `work_order_id`, SK: `comment_id`)
Backfill stats (initial run): 14,825 POs scanned → 9,900 with extractable site codes → 1,100 unique sites.
## Architecture
**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`.
- **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`.
## CI/CD
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)
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.
Branch protection on `main` — all changes through PR.
**Branch protection:** `main` requires a PR (no direct push), no deletion, no force push.
## Setup
1. Bootstrap CDK: `cdk bootstrap aws://{AccountId}/us-east-1`
2. Store Anthropic API keys:
1. Bootstrap CDK in the account if you haven't already: `cdk bootstrap aws://{AccountId}/us-east-1`.
2. Store the Anthropic API key:
```bash
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-..."
aws secretsmanager create-secret \
--name po-ingest/anthropic-api-key \
--secret-string "sk-ant-..."
```
3. Deploy both stacks:
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:
```bash
cd cdk
pip install -r requirements.txt
cdk deploy --all
cdk deploy
```
4. CloudFormation outputs include `WebUIUrl` for each stack's dashboard.
5. The `WebUIUrl` CloudFormation output is the dashboard URL.
## Scripts
## Reprocessing
To re-run the processor against every email still sitting in `inbound/` (useful after a parser change):
**Reprocess PO emails** (re-run parser against all emails still in S3):
```bash
python scripts/reprocess.py # dry-run
python scripts/reprocess.py --execute # invoke po-email-processor for each
python scripts/reprocess.py # dry-run — lists keys
python scripts/reprocess.py --execute # invokes po-email-processor for each
```
**Backfill verified sites** (one-time scan of historical POs):
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):
```bash
python scripts/backfill_sites.py
```
## 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
```
Uses the same extraction logic as the Lambda. Idempotent — safe to re-run.

15
buildspec.yml Normal file
View file

@ -0,0 +1,15 @@
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,22 +1,12 @@
#!/usr/bin/env python3
import aws_cdk as cdk
from po_stack import PoIngestStack
from wo_stack import WorkorderIngestStack
from stack import PoIngestStack
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/po/email_processor",
"../lambdas/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/po/web_ui"),
code=lambda_.Code.from_asset("../lambdas/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/po/site_extractor"),
code=lambda_.Code.from_asset("../lambdas/site_extractor"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,

View file

@ -1,185 +0,0 @@
"""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,10 +5,8 @@ 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
@ -33,7 +31,7 @@ def render_badge(value, color_map):
if not value:
value = "unknown"
color = color_map.get(value.lower(), "#9ca3af")
label = esc(value.replace("_", " ").title())
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>'
@ -64,51 +62,50 @@ def fmt_currency(val):
def render_po_detail(po):
po_number = esc(po.get("po_number", ""))
po_number = 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", 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),
("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")),
]
ship_to = po.get("ship_to") or {}
if any(ship_to.values()):
ship_parts = []
if ship_to.get("name"):
ship_parts.append(esc(ship_to["name"]))
ship_parts.append(ship_to["name"])
if ship_to.get("address"):
ship_parts.append(esc(ship_to["address"]))
ship_parts.append(ship_to["address"])
if ship_to.get("location_code"):
ship_parts.append(f"Location: {esc(ship_to['location_code'])}")
ship_parts.append(f"Location: {ship_to['location_code']}")
if ship_to.get("attn"):
ship_parts.append(f"Attn: {esc(ship_to['attn'])}")
ship_parts.append(f"Attn: {ship_to['attn']}")
fields.append(("Ship To", "<br>".join(ship_parts)))
view_url = po.get("view_order_url")
if view_url and view_url.startswith(("https://", "http://")):
escaped_url = esc(view_url, quote=True)
if view_url:
fields.append(
(
"Coupa Link",
f'<a href="{escaped_url}" target="_blank" style="color:#3b82f6;">View in Coupa</a>',
f'<a href="{view_url}" target="_blank" style="color:#3b82f6;">View in Coupa</a>',
)
)
@ -127,17 +124,17 @@ def render_po_detail(po):
if line_items:
rows = ""
for item in line_items:
qty = esc(str(item.get("quantity", "") or ""))
unit = esc(str(item.get("unit", "") or ""))
price = esc(str(item.get("price", "") or ""))
qty = item.get("quantity", "") or ""
unit = item.get("unit", "") or ""
price = item.get("price", "") or ""
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;">
<td style="padding:10px;font-size:14px;">{esc(str(item.get("description", "")))}</td>
<td style="padding:10px;font-size:14px;">{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;">{esc(str(item.get("need_by", "") or ""))}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{item.get("need_by", "") or ""}</td>
</tr>"""
items_html = f"""
@ -187,16 +184,16 @@ def render_po_detail(po):
def render_po_list(purchase_orders):
rows = ""
for po in purchase_orders:
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", ""))
po_number = po.get("po_number", "")
supplier = (po.get("supplier") or {}).get("name", "")
site_code = po.get("site_code", "")
trade = po.get("trade", "")
status = po.get("po_status", "")
total = fmt_currency(po.get("total_amount"))
processed = esc((po.get("processed_at") or "")[:16])
processed = (po.get("processed_at") or "")[:16]
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location={esc(json.dumps(f"/po?id={po.get("po_number", "")}"), quote=True)}">
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location='/po?id={po_number}'">
<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

@ -1,271 +0,0 @@
"""
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

@ -1,2 +0,0 @@
anthropic>=0.101.0
boto3>=1.43.6

View file

@ -1,256 +0,0 @@
"""
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

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

View file

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

View file

@ -1,47 +0,0 @@
#!/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()