diff --git a/.github/dependabot.yml b/.github/dependabot.yml index c571f0e..9717d1b 100644 --- a/.github/dependabot.yml +++ b/.github/dependabot.yml @@ -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: diff --git a/.gitignore b/.gitignore index 5d81fa3..1dafb99 100644 --- a/.gitignore +++ b/.gitignore @@ -9,4 +9,4 @@ node_modules/ cdk.out/ .env *.eml -lambdas/*/package/ +lambdas/*/*/package/ diff --git a/README.md b/README.md index ae75e9d..16f74f9 100644 --- a/README.md +++ b/README.md @@ -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 +``` diff --git a/buildspec.yml b/buildspec.yml deleted file mode 100644 index da69dde..0000000 --- a/buildspec.yml +++ /dev/null @@ -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 diff --git a/cdk/app.py b/cdk/app.py index 2fac78d..6f237be 100644 --- a/cdk/app.py +++ b/cdk/app.py @@ -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() diff --git a/cdk/stack.py b/cdk/po_stack.py similarity index 97% rename from cdk/stack.py rename to cdk/po_stack.py index 6f4bc86..f9d0d6c 100644 --- a/cdk/stack.py +++ b/cdk/po_stack.py @@ -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, diff --git a/cdk/wo_stack.py b/cdk/wo_stack.py new file mode 100644 index 0000000..dd29781 --- /dev/null +++ b/cdk/wo_stack.py @@ -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" + ) diff --git a/lambdas/email_processor/handler.py b/lambdas/po/email_processor/handler.py similarity index 100% rename from lambdas/email_processor/handler.py rename to lambdas/po/email_processor/handler.py diff --git a/lambdas/email_processor/requirements.txt b/lambdas/po/email_processor/requirements.txt similarity index 100% rename from lambdas/email_processor/requirements.txt rename to lambdas/po/email_processor/requirements.txt diff --git a/lambdas/site_extractor/handler.py b/lambdas/po/site_extractor/handler.py similarity index 100% rename from lambdas/site_extractor/handler.py rename to lambdas/po/site_extractor/handler.py diff --git a/lambdas/web_ui/handler.py b/lambdas/po/web_ui/handler.py similarity index 79% rename from lambdas/web_ui/handler.py rename to lambdas/po/web_ui/handler.py index 90ef3c6..beaa174 100644 --- a/lambdas/web_ui/handler.py +++ b/lambdas/po/web_ui/handler.py @@ -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'{label}' @@ -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", "
".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'View in Coupa', + f'View in Coupa', ) ) @@ -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""" - {item.get("description", "")} + {esc(str(item.get("description", "")))} {qty} {unit} {price} {fmt_currency(item.get("amount"))} - {item.get("need_by", "") or ""} + {esc(str(item.get("need_by", "") or ""))} """ 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""" - + {po_number} {supplier} {site_code} diff --git a/lambdas/wo/email_processor/__init__.py b/lambdas/wo/email_processor/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/lambdas/wo/email_processor/handler.py b/lambdas/wo/email_processor/handler.py new file mode 100644 index 0000000..2abdc8f --- /dev/null +++ b/lambdas/wo/email_processor/handler.py @@ -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"} diff --git a/lambdas/wo/email_processor/requirements.txt b/lambdas/wo/email_processor/requirements.txt new file mode 100644 index 0000000..acef34e --- /dev/null +++ b/lambdas/wo/email_processor/requirements.txt @@ -0,0 +1,2 @@ +anthropic>=0.42.0 +boto3>=1.35.0 diff --git a/lambdas/wo/web_ui/__init__.py b/lambdas/wo/web_ui/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/lambdas/wo/web_ui/handler.py b/lambdas/wo/web_ui/handler.py new file mode 100644 index 0000000..4af9120 --- /dev/null +++ b/lambdas/wo/web_ui/handler.py @@ -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'{label}' + + +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"{commenter} — " 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""" +
+
+ {badge} {commenter_str}{created} +
+ {"
" + text + "
" if text else ""} +
""" + else: + events_html = '

No events yet.

' + + 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""" +
+
{label}
+
{value}
+
""" + + return f""" + + + + + WO {wo_id} - Sea Haven + + + +
+
+ ← All Work Orders +
+
+

Work Order {wo_id}

+ {details_html} +
+
+

Events ({len(events)})

+ {events_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""" + + {wo_id} + {desc} + {site} + {render_badge(record_type, RECORD_TYPE_COLORS)} + {render_badge(status, STATUS_COLORS)} + {due} + {updated} + """ + + return f""" + + + + + Work Orders - Sea Haven + + + +
+
+

Work Orders

+ {len(work_orders)} total +
+
+ + + + + + + + + + + + + + {rows if rows else ''} + +
WO #DescriptionSiteLast ActionStatusDue DateUpdated
No work orders yet.
+
+
+ +""" + + +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": "

Work order not found

", + } + 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, + } diff --git a/lambdas/wo/web_ui/requirements.txt b/lambdas/wo/web_ui/requirements.txt new file mode 100644 index 0000000..3f3a438 --- /dev/null +++ b/lambdas/wo/web_ui/requirements.txt @@ -0,0 +1 @@ +boto3>=1.35.0 diff --git a/scripts/backfill_sites.py b/scripts/backfill_sites.py index 121d5ca..f2b9909 100644 --- a/scripts/backfill_sites.py +++ b/scripts/backfill_sites.py @@ -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 diff --git a/test_local.py b/test_local.py new file mode 100644 index 0000000..77e8010 --- /dev/null +++ b/test_local.py @@ -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()