procurement-ingest/lambdas/po/email_processor/handler.py
Adam Moussa 5785f39590 Merge PO revisions and handle out-of-order events
save_revision did a full put_item overwrite, so a revision omitting
line_items/supplier permanently deleted them. save_new_po used a
conditional put that silently dropped the PO when an out-of-order
cancellation had already created a skeleton row.

Switch both to field-level merge update_items: a revision now SETs
only the fields it carries, and a new_po backfills data into a
pre-existing Cancelled skeleton while preserving the Cancelled
status. No email can now delete data established by an earlier one.
2026-06-17 11:37:35 -04:00

541 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
PO email processor Lambda.
Triggered by S3 events when SES delivers a Coupa PO email.
Parses the raw email, sends it to Claude for structured extraction,
then writes the result to the purchase-orders DynamoDB table.
"""
import email
import email.utils
import json
import logging
import os
import re
from datetime import datetime
from decimal import Decimal
from email import policy
import anthropic
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
s3 = boto3.client("s3")
dynamodb = boto3.resource("dynamodb")
PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders")
ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN")
# Allowlist of sender domains permitted to create/modify POs. Coupa sends Amazon
# procurement notifications from the coupahost.com family; the leading dot means
# "this domain or any subdomain". Override via the ALLOWED_SENDER_DOMAINS env var
# (comma-separated) without a redeploy of code. An attacker who can email
# amazon_po@int.seahaven.com but cannot forge a verified sender in this list is
# rejected before any DynamoDB write.
DEFAULT_ALLOWED_SENDER_DOMAINS = "coupahost.com,amazon.com"
ALLOWED_SENDER_DOMAINS = [
d.strip().lower().lstrip("@")
for d in os.environ.get(
"ALLOWED_SENDER_DOMAINS", DEFAULT_ALLOWED_SENDER_DOMAINS
).split(",")
if d.strip()
]
EXTRACTION_PROMPT = """\
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" | "revision" | "cancellation",
"po_number": "string or null",
"po_status": "string or null",
"source_system": "coupa",
"submitted_by": "string or null",
"on_behalf_of": "string or null",
"order_date": "string or null",
"revision_date": "string or null",
"last_opened": "string or null",
"acknowledged_at": "string or null",
"payment_terms": "string or null",
"requisition_number": "string or null",
"department": "string or null",
"view_order_url": "URL string or null",
"supplier": {
"name": "string or null"
},
"site_code": "string or null",
"ship_to": {
"name": "string or null",
"address": "string or null",
"street": "string or null",
"city": "string or null",
"state": "string or null",
"zip": "string or null",
"location_code": "string or null",
"attn": "string or null"
},
"total_amount": 0.0,
"currency": "USD",
"fiscal_year": "string or null",
"trade": "string or null",
"coupa_category": "string or null",
"line_items": [
{
"description": "string",
"amount": 0.0,
"currency": "USD",
"need_by": "date string or null",
"category": "string or null",
"account_code": "string or null",
"period": "string or null",
"quantity": "string or null",
"unit": "string or null",
"price": "string or null"
}
]
}
## 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
"""
def get_anthropic_client() -> anthropic.Anthropic:
"""Create Anthropic client, fetching API key from Secrets Manager.
Requires ANTHROPIC_API_KEY_SECRET_ARN to be set. The previous silent
fallback to a plaintext ANTHROPIC_API_KEY env var is removed: a misconfigured
deploy must fail loudly rather than quietly run on an unmanaged key.
"""
if not ANTHROPIC_API_KEY_SECRET_ARN:
raise RuntimeError(
"ANTHROPIC_API_KEY_SECRET_ARN is not set; refusing to fall back to a "
"plaintext API key. Configure the Secrets Manager 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)
def extract_sender_domain(sender: str) -> str | None:
"""Extract the lowercased domain from a From header value.
Handles "Name <user@domain>" and bare "user@domain" forms.
"""
if not sender:
return None
_, addr = email.utils.parseaddr(sender)
if "@" not in addr:
return None
return addr.rsplit("@", 1)[1].strip().lower()
def is_sender_allowed(sender: str) -> bool:
"""Return True if the sender domain is in the configured allowlist.
A domain matches if it equals an allowlist entry or is a subdomain of one
(e.g. "notifications.coupahost.com" matches "coupahost.com").
"""
domain = extract_sender_domain(sender)
if not domain:
return False
for allowed in ALLOWED_SENDER_DOMAINS:
if domain == allowed or domain.endswith("." + allowed):
return True
return False
def ses_auth_failed(email_data: dict) -> bool:
"""Return True if SES recorded a hard SPF or DKIM failure for this email.
SES (when receipt-rule spam/virus/auth scanning is enabled) stamps the stored
object with X-SES-Spam-Verdict / X-SES-Virus-Verdict and an
Authentication-Results header carrying spf=/dkim= results. We only block on an
explicit "fail" so that mails delivered before scanning is enabled (no header)
are not silently dropped; the domain allowlist remains the primary gate.
"""
spam = (email_data.get("ses_spam_verdict") or "").upper()
virus = (email_data.get("ses_virus_verdict") or "").upper()
if spam == "FAIL" or virus == "FAIL":
return True
auth = (email_data.get("authentication_results") or "").lower()
if "spf=fail" in auth or "dkim=fail" in auth:
return True
return False
def parse_raw_email(raw_bytes: bytes) -> dict:
"""Parse a raw MIME email into subject, sender, and body text."""
msg = email.message_from_bytes(raw_bytes, policy=policy.default)
subject = msg.get("Subject", "")
sender = msg.get("From", "")
to = msg.get("To", "")
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,
"date": date,
"body": body,
# SES stamps these on the stored object when receipt-rule scanning is on.
"ses_spam_verdict": msg.get("X-SES-Spam-Verdict", ""),
"ses_virus_verdict": msg.get("X-SES-Virus-Verdict", ""),
"authentication_results": msg.get("Authentication-Results", ""),
}
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"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=2048,
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(), parse_float=Decimal)
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
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
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 _merge_update(po_number: str, fields: dict):
"""Apply a merge (SET-only) update of the given fields onto a PO record.
Only the fields supplied are written; absent fields are left untouched, so a
partial payload can never delete data that an earlier email established. The
record is created if it does not exist (DynamoDB update_item upsert).
"""
table = dynamodb.Table(PO_TABLE)
set_parts = []
attr_names = {}
attr_values = {}
for key, value in fields.items():
if value is None or key == "po_number":
continue
name_ph = f"#{key}"
val_ph = f":{key}"
attr_names[name_ph] = key
attr_values[val_ph] = value
set_parts.append(f"{name_ph} = {val_ph}")
if not set_parts:
return
table.update_item(
Key={"po_number": po_number},
UpdateExpression="SET " + ", ".join(set_parts),
ExpressionAttributeNames=attr_names,
ExpressionAttributeValues=attr_values,
)
def _is_cancelled(po_number: str) -> bool:
"""Return True if the PO already exists with a Cancelled status."""
table = dynamodb.Table(PO_TABLE)
existing = table.get_item(Key={"po_number": po_number}).get("Item")
return bool(existing) and existing.get("po_status") == "Cancelled"
def save_new_po(parsed: dict):
"""Create a PO, merging into any pre-existing record.
Uses a merge update rather than a conditional put so that an out-of-order
cancellation (which leaves a Cancelled skeleton) is filled in with the full
PO data instead of the new_po being silently dropped. The cancellation marker
(po_status=Cancelled, cancelled_at) is preserved; new_po data backfills the
remaining fields.
"""
po_number = parsed["po_number"]
fields = {k: v for k, v in parsed.items() if v is not None}
if _is_cancelled(po_number):
# Preserve the cancellation: don't overwrite po_status/cancelled_at.
fields.pop("po_status", None)
fields.pop("cancelled_at", None)
logger.info(
f"PO {po_number} was cancelled before new_po arrived; "
f"backfilling data and keeping Cancelled status"
)
_merge_update(po_number, fields)
logger.info(f"Created/merged PO {po_number}")
def save_revision(parsed: dict):
"""Merge revised data into an existing PO without deleting omitted fields.
A revision email often omits unchanged sections (line_items, supplier). The
previous full-overwrite put_item permanently dropped those. This SETs only the
fields present in the revision, leaving everything else intact.
"""
po_number = parsed["po_number"]
fields = {k: v for k, v in parsed.items() if v is not None}
if _is_cancelled(po_number):
# A revision must not silently un-cancel a PO.
fields.pop("po_status", None)
_merge_update(po_number, fields)
logger.info(f"Revised PO {po_number}")
def save_cancellation(parsed: dict):
"""Mark a PO Cancelled, creating a minimal skeleton if it doesn't exist yet.
If the cancellation arrives before the new_po, the skeleton it creates is
later backfilled by save_new_po (which preserves this Cancelled status), so no
PO data is lost on out-of-order delivery.
"""
table = dynamodb.Table(PO_TABLE)
table.update_item(
Key={"po_number": parsed["po_number"]},
UpdateExpression="SET po_status = :status, cancelled_at = :cancelled_at, raw_s3_key = :s3_key",
ExpressionAttributeValues={
":status": "Cancelled",
":cancelled_at": parsed.get("processed_at", datetime.utcnow().isoformat()),
":s3_key": parsed.get("raw_s3_key", ""),
},
)
logger.info(f"Cancelled PO {parsed['po_number']}")
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"]
s3_key = f"s3://{bucket}/{key}"
logger.info(f"Processing email: {s3_key}")
response = s3.get_object(Bucket=bucket, Key=key)
raw_email = response["Body"].read()
email_data = parse_raw_email(raw_email)
logger.info(f"Subject: {email_data['subject']}")
# Authz: only act on email from an allowlisted sender domain that passed
# SES auth checks. Anyone can email amazon_po@int.seahaven.com, but only
# legitimate Coupa/Amazon senders may create or mutate POs.
sender = email_data.get("sender", "")
if not is_sender_allowed(sender):
logger.warning(
f"Rejecting email from disallowed sender '{sender}' "
f"(domain not in allowlist): {key}"
)
continue
if ses_auth_failed(email_data):
logger.warning(
f"Rejecting email from '{sender}' due to SES SPF/DKIM/spam "
f"failure: {key}"
)
continue
parsed = extract_with_claude(email_data)
logger.info(
f"Parsed: type={parsed.get('email_type')}, po={parsed.get('po_number')}"
)
if not parsed.get("po_number"):
logger.warning(f"No PO number found in email, skipping: {key}")
continue
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)
return {"statusCode": 200, "body": "OK"}