Merge pull request #4 from Sea-Haven-Industries/feature/align-po-schema

Align PO schema with enriched records and improve extraction prompt
This commit is contained in:
Adam Moussa 2026-05-01 19:58:28 -04:00 • committed by GitHub
commit f942755495
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 216 additions and 43 deletions

View file

@ -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:

View file

@ -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 ---

View file

@ -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"}

View file

@ -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"""
<tr style="border-bottom:1px solid #f1f5f9;">
<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;">{item.get('need_by', '') or ''}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{item.get('category', '') or ''}</td>
</tr>"""
items_html = f"""
@ -128,9 +140,11 @@ def render_po_detail(po):
<thead>
<tr>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Description</th>
<th style="text-align:right;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Qty</th>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Unit</th>
<th style="text-align:right;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Price</th>
<th style="text-align:right;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Amount</th>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Need By</th>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Category</th>
</tr>
</thead>
<tbody>{rows}</tbody>
@ -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):
<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;">{render_badge(email_type, EMAIL_TYPE_COLORS)}</td>
<td style="padding:12px;font-weight:500;">{site_code}</td>
<td style="padding:12px;font-size:13px;">{trade}</td>
<td style="padding:12px;">{render_badge(status, STATUS_COLORS)}</td>
<td style="padding:12px;text-align:right;">{total}</td>
<td style="padding:12px;color:#64748b;font-size:13px;">{processed}</td>
@ -209,14 +225,15 @@ def render_po_list(purchase_orders):
<tr>
<th>PO #</th>
<th>Supplier</th>
<th>Type</th>
<th>Site</th>
<th>Trade</th>
<th>Status</th>
<th style="text-align:right;">Amount</th>
<th>Processed</th>
</tr>
</thead>
<tbody>
{rows if rows else '<tr><td colspan="6" style="padding:40px;text-align:center;color:#94a3b8;">No purchase orders yet.</td></tr>'}
{rows if rows else '<tr><td colspan="7" style="padding:40px;text-align:center;color:#94a3b8;">No purchase orders yet.</td></tr>'}
</tbody>
</table>
</div>