diff --git a/README.md b/README.md index 3e67e7d..ae75e9d 100644 --- a/README.md +++ b/README.md @@ -9,13 +9,26 @@ Coupa purchase-order email ingestion pipeline. SES receives Amazon PO emails, Cl 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_AND_OLD_IMAGES) on `purchase-orders` feeds two downstream consumers: +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). +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. + A separate `po-web-ui` Lambda (Function URL, unauthenticated) renders a simple HTML dashboard scanning the table. +### PO record schema + +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` + ### Verified-sites pipeline The `po-ingest-site-extractor` Lambda is triggered by the DynamoDB Stream on every PO INSERT/MODIFY. It: diff --git a/cdk/stack.py b/cdk/stack.py index 21ec4c5..92cd7c0 100644 --- a/cdk/stack.py +++ b/cdk/stack.py @@ -45,7 +45,7 @@ class PoIngestStack(Stack): ), billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, removal_policy=RemovalPolicy.RETAIN, - stream=dynamodb.StreamViewType.NEW_AND_OLD_IMAGES, + stream=dynamodb.StreamViewType.NEW_IMAGE, ) # --- Secrets Manager for Anthropic API key --- diff --git a/lambdas/email_processor/handler.py b/lambdas/email_processor/handler.py index eea5a8e..51daa88 100644 --- a/lambdas/email_processor/handler.py +++ b/lambdas/email_processor/handler.py @@ -28,14 +28,14 @@ PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders") ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN") EXTRACTION_PROMPT = """\ -You are an email parser for a purchase order system. -The emails come from Coupa (a procurement platform) and contain purchase order -notifications from Amazon. +You are an email parser for a purchase order ingest pipeline. +The emails are Coupa procurement platform notifications containing purchase order +data from Amazon. Analyze the following email and extract structured data. Return ONLY valid JSON with these fields: { - "email_type": "new_po" | "cancellation", + "email_type": "new_po" | "revision" | "cancellation", "po_number": "string or null", "po_status": "string or null", "source_system": "coupa", @@ -65,6 +65,9 @@ Analyze the following email and extract structured data. Return ONLY valid JSON }, "total_amount": 0.0, "currency": "USD", + "fiscal_year": "string or null", + "trade": "string or null", + "coupa_category": "string or null", "line_items": [ { "description": "string", @@ -73,25 +76,134 @@ Analyze the following email and extract structured data. Return ONLY valid JSON "need_by": "date string or null", "category": "string or null", "account_code": "string or null", - "period": "string or null" + "period": "string or null", + "quantity": "string or null", + "unit": "string or null", + "price": "string or null" } ] } -Rules: -- "email_type" detection: - - "new_po": email announces a new or revised purchase order - - "cancellation": email announces a PO has been cancelled -- Extract the PO number from the email (e.g., "2D-18206023") -- Extract all line items with their descriptions, amounts, and metadata -- Ship-to address should include the full address, location code, and attention line -- Ship-to street, city, state, and zip should be parsed from the address into separate fields -- "site_code" is the Amazon facility code (e.g., "SNY5", "DFW6", "WND1") — a 3-5 character alphanumeric code identifying the delivery site. Look for it in: - - The ship-to name, e.g., "Amazon.com Services LLC - SNY5" or "Amazon.com Services LLC (WFB1)" - - The ATTN line, e.g., "ATTN: Wagon Wheel DS - WTN1" - - Line item descriptions, e.g., "WND1 - 2024 - Plumbing PM" - - Anywhere else in the email where a facility code appears - - If the ship-to name IS the site code (e.g., just "DBU2"), use that +## email_type detection + +- "new_po": email announces a new purchase order being issued +- "revision": email announces a revised/updated purchase order (look for "revised" in subject or body) +- "cancellation": email announces a PO has been cancelled + +## PO number + +Extract from the email subject or body. Format is a prefix + hyphen + digits: +- "2D-18206023", "FK-21088051", "B187-17955555" + +## site_code extraction + +The site code is the Amazon facility code — a 3-5 character alphanumeric code identifying +the delivery site. Check these locations in order: + +1. Ship-to name in parentheses: "Amazon.com Services LLC (KLAL)" → KLAL +2. Ship-to name after dash: "Amazon.com Services LLC - SNY5" → SNY5 +3. Ship-to ATTN line with dash or en-dash: "ATTN: Wagon Wheel DS Station –WKY3" → WKY3 +4. Ship-to ATTN line directly: "Attn: HJX1" → HJX1 +5. Ship-to name IS the code: if the name is just "DBU2" or similar, use it +6. Line item description prefix: "DYO1 - Sea Haven Ind - Plumbing Repairs" → DYO1 +7. Line item description in brackets: "[HMK4] Assemble 3 Wire Security Cages" → HMK4 + +**Not site codes — do not extract these as site_code:** +- RME (Amazon Reliability Maintenance Engineering department) +- BBM (Coupa description format tag) +- JLL (Jones Lang LaSalle — facilities management vendor) +- PARAG, ERIK (vendor/person names) +- Industry acronyms: HVAC, LED, PVC, ADA, OSHA, EMR, BMS, DDC, MRO, NTE, EST + +If the only candidate matches this skip list, set site_code to null. + +## Ship-to address parsing + +Parse the full address into separate fields. Be aware of these common issues: +- State abbreviation may be missing entirely (e.g., "Tucson, 85704" with no state) +- Zip codes may lack leading zeros (e.g., "MA 2149" should be zip "02149", "NJ 7001" should be "07001") +- City names may be misspelled (e.g., "Charoltte" for Charlotte) — extract as-is, do not correct +- Format varies: "City, ST - ZIP", "City, ST ZIP", "City, ZIP" (no state) + +If state cannot be determined from the address, set ship_to.state to null. + +## fiscal_year + +The calendar year the work covers. Determine from: +1. The order_date year (primary source) +2. Need-by dates on line items +3. Year in line item descriptions (e.g., "HVB2 - 2025 - Plumbing PM" → "2025") + +Use the 4-digit year string (e.g., "2025"). + +## trade classification + +Classify the primary trade from line item descriptions. Use the FIRST match in priority order: + +**Plumbing - PM**: "plumbing pm", "plumbing preventative", "plumbing maintenance", + or BBM format: "Plumbing - Backflow", "Plumbing - Water Heater - Install/Repair" + +**Plumbing - Reactive**: "plumbing" with: "reactive", "emergency", "repair", "clog", + "unclog", "leak", "flood", "sewer", "drain", "grease trap", "jetter", "water line", + "toilet", "faucet", "urinal", "pipe" + +**Electrical**: "electrical", "lighting", "ballast", "outlet", "circuit", "panel", + "generator", "transformer", "conduit" (but NOT if "dock door" context) + +**HVAC**: "hvac", "heating", "cooling", "air conditioning", "RTU", "AHU", "VAV", + "refrigerant", "thermostat", "ductwork" + +**Dock Doors**: "dock door", "dock leveler", "dock plate", "dock seal", "dock bumper" + +**Doors**: "door", "overhead door", "roll-up", "automatic door", "access door" + (only if not matched by Dock Doors above) + +**Signage**: "sign", "banner", "wayfinding", "marquee", "directional" + +**Carpentry**: "carpentry", "cabinet", "millwork", "trim", "shelving", "framing" + +**Fencing/Gates**: "fence", "fencing", "gate", "bollard" (not "dock gate") + +**Conveyance/MHE**: "conveyor", "MHE", "material handling", "sortation" + +**Painting**: "paint", "painting", "primer", "coating", "touch-up" + +**Flooring**: "floor", "tile", "carpet", "epoxy", "polishing" + +**Janitorial**: "janitorial", "cleaning", "custodial", "pressure wash", "power wash" + +**Fire/Life Safety**: "fire", "sprinkler", "extinguisher", "fire alarm", "suppression" + +**Landscaping/Yard**: "landscape", "lawn", "tree", "yard", "mowing", "irrigation" + +**Roofing**: "roof", "roofing", "gutter", "downspout" + +**Security/Locksmith**: "lock", "key", "access control", "camera", "security", "CCTV" + +**Snow Removal**: "snow", "ice", "salt", "de-ice", "plow" + +**PO Uplift**: description is exactly or primarily "PO Uplift" + +**General Building - Emergency**: "EMER" prefix, or "emergency" in a general building context + +**General Building - Handyman**: BBM format "General Building - General Building Technician" + +**General Building - Project**: BBM format "General Building - General Building Project" + +**General Building**: any remaining facility maintenance work + +If a PO has multiple line items with different trades, set "trade" to the primary +(non-uplift, non-materials) trade. If genuinely mixed, use the trade of the highest-value line item. + +## coupa_category + +The Coupa commodity/category field if present in the email (e.g., "Maintenance - Facilities", +"Plumbing Equipment & Materials"). This is Coupa's own classification, not the trade field. + +## General rules + +- Extract all line items with descriptions, amounts, and metadata +- "quantity", "unit" (e.g., "EACH", "HR"), and "price" (unit price) should be extracted when present - total_amount should be the numeric total in USD - 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 @@ -173,16 +285,37 @@ def extract_with_claude(email_data: dict) -> dict: return json.loads(response_text.strip(), parse_float=Decimal) -def save_new_po(parsed: dict, s3_key: str): - """Insert a new PO into DynamoDB. Skips if po_number already exists.""" - table = dynamodb.Table(PO_TABLE) - now = datetime.utcnow().isoformat() +def pad_zip(zip_code: str | None) -> str | None: + if not zip_code: + return zip_code + clean = zip_code.strip().split("-")[0] + if clean.isdigit() and len(clean) < 5: + return clean.zfill(5) + zip_code.strip()[len(clean):] + return zip_code - # Add metadata fields + +def enrich_parsed(parsed: dict, s3_key: str, email_subject: str): + """Add metadata and promote nested fields to top level.""" + now = datetime.utcnow().isoformat() parsed["raw_s3_key"] = s3_key parsed["processed_at"] = now + parsed["data_source"] = "email" + parsed["email_subject"] = email_subject - # Build item, stripping None values + ship_to = parsed.get("ship_to") or {} + if ship_to.get("address"): + parsed["ship_to_raw"] = ship_to["address"] + if ship_to.get("state"): + parsed["state"] = ship_to["state"] + + if ship_to.get("zip"): + ship_to["zip"] = pad_zip(ship_to["zip"]) + return parsed + + +def save_new_po(parsed: dict): + """Insert a new PO into DynamoDB. Skips if po_number already exists.""" + table = dynamodb.Table(PO_TABLE) item = {k: v for k, v in parsed.items() if v is not None} try: @@ -195,7 +328,16 @@ def save_new_po(parsed: dict, s3_key: str): logger.info(f"PO {parsed['po_number']} already exists, skipping insert") -def save_cancellation(parsed: dict, s3_key: str): +def save_revision(parsed: dict): + """Update an existing PO with revised data, or insert if it doesn't exist yet.""" + table = dynamodb.Table(PO_TABLE) + item = {k: v for k, v in parsed.items() if v is not None} + + table.put_item(Item=item) + logger.info(f"Revised PO {parsed['po_number']}") + + +def save_cancellation(parsed: dict): """Update an existing PO's status to Cancelled.""" table = dynamodb.Table(PO_TABLE) @@ -204,8 +346,8 @@ def save_cancellation(parsed: dict, s3_key: str): UpdateExpression="SET po_status = :status, cancelled_at = :cancelled_at, raw_s3_key = :s3_key", ExpressionAttributeValues={ ":status": "Cancelled", - ":cancelled_at": datetime.utcnow().isoformat(), - ":s3_key": s3_key, + ":cancelled_at": parsed.get("processed_at", datetime.utcnow().isoformat()), + ":s3_key": parsed.get("raw_s3_key", ""), }, ) logger.info(f"Cancelled PO {parsed['po_number']}") @@ -220,15 +362,12 @@ def handler(event, context): logger.info(f"Processing email: {s3_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')}, po={parsed.get('po_number')}") @@ -236,10 +375,14 @@ def handler(event, context): logger.warning(f"No PO number found in email, skipping: {key}") continue - # Route by email type - if parsed.get("email_type") == "cancellation": - save_cancellation(parsed, s3_key) + parsed = enrich_parsed(parsed, s3_key, email_data["subject"]) + + email_type = parsed.get("email_type") + if email_type == "cancellation": + save_cancellation(parsed) + elif email_type == "revision": + save_revision(parsed) else: - save_new_po(parsed, s3_key) + save_new_po(parsed) return {"statusCode": 200, "body": "OK"} diff --git a/lambdas/web_ui/handler.py b/lambdas/web_ui/handler.py index 22d6d44..fecf9e7 100644 --- a/lambdas/web_ui/handler.py +++ b/lambdas/web_ui/handler.py @@ -47,6 +47,7 @@ STATUS_COLORS = { EMAIL_TYPE_COLORS = { "new_po": "#3b82f6", + "revision": "#f59e0b", "cancellation": "#ef4444", } @@ -71,6 +72,11 @@ def render_po_detail(po): ("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")), @@ -78,6 +84,7 @@ def render_po_detail(po): ("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")), ] @@ -113,12 +120,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 '' rows += f""" {item.get('description', '')} + {qty} + {unit} + {price} {fmt_currency(item.get('amount'))} {item.get('need_by', '') or ''} - {item.get('category', '') or ''} """ items_html = f""" @@ -128,9 +140,11 @@ def render_po_detail(po): Description + Qty + Unit + Price Amount Need By - Category {rows} @@ -168,8 +182,9 @@ def render_po_list(purchase_orders): 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", "") status = po.get("po_status", "") - email_type = po.get("email_type", "") total = fmt_currency(po.get("total_amount")) processed = (po.get("processed_at") or "")[:16] @@ -177,7 +192,8 @@ def render_po_list(purchase_orders): {po_number} {supplier} - {render_badge(email_type, EMAIL_TYPE_COLORS)} + {site_code} + {trade} {render_badge(status, STATUS_COLORS)} {total} {processed} @@ -209,14 +225,15 @@ def render_po_list(purchase_orders): PO # Supplier - Type + Site + Trade Status Amount Processed - {rows if rows else 'No purchase orders yet.'} + {rows if rows else 'No purchase orders yet.'}