mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 11:53:13 +00:00
Some checks are pending
Deploy / deploy (push) Waiting to run
* Add deterministic template parser for WO emails The workorder-email-processor sends every one of ~22.9k emails/month to an LLM, but ~93.6% are the plain-text "AMAZON UPDATE WO DETAILS" comment template and ~6.4% the HTML "AMAZON assign Work Order" template. Parse those two shapes deterministically, offline, so the AI call is reserved for the long tail. The module is pure (no boto3, no network). try_deterministic_parse classifies by subject, extracts the shared contract fields, and returns a result ONLY when it passes a strict fail-closed validation gate: exact contract-key set, subject/id agreement, the literal "Work Order: <id>" double space, per-type required fields, site-code shape, and a label-bleed guard so a value that over-ran into the next field fails. Any miss, drift, or extractor exception yields None so the caller falls back to the AI extractor -- data is never corrupted, only the fallback rate rises. Refs: #23 * Migrate WO processor to Bedrock and fix comment_id collision Switch the AI path from the Anthropic SDK to bedrock-runtime InvokeModel on the inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0 (BEDROCK_MODEL_ID env), so parsing no longer needs a provider API key or Secrets Manager secret. The EXTRACTION_PROMPT and JSON contract are kept byte-identical, so the AI-fallback output is unchanged. Try the new deterministic template parser first and only call Bedrock on a miss/invalid result. Fix issue #23: the WorkOrderComments range key was work_order_id#<comment_time>, so two emails on one WO with an identical or absent comment time collided and overwrote each other. Derive a 12-hex suffix from the S3 object key alone -- deterministic, so an async retry of the same object is byte-identical (idempotent) while distinct emails get distinct keys -- and keep wall-clock now() out of the key (literal 'nocomment' segment when comment_time is absent). Also emit one CloudWatch EMF line per record (Seahaven/WorkorderIngest ParseOutcome, dimensioned by ParseMethod/TemplateId) for parse-outcome observability, replace the deprecated datetime.utcnow() with datetime.now(timezone.utc), and drop the anthropic dependency. Refs: #23 * Migrate PO processor to Bedrock Switch the PO email processor's AI extraction from the Anthropic SDK to bedrock-runtime InvokeModel on the inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0 (BEDROCK_MODEL_ID env), so it no longer needs a provider API key or Secrets Manager secret. PO parsing stays fully AI -- only the provider changes. The EXTRACTION_PROMPT is kept byte-identical and the Bedrock text output is still decoded with json.loads(..., parse_float=Decimal), which DynamoDB requires (it rejects floats). Replace the deprecated datetime.utcnow() with datetime.now(timezone.utc) and drop the anthropic dependency. * Grant Bedrock IAM, drop Anthropic secrets, add fallback alarm Both stacks moved their processors from the Anthropic API to the Bedrock inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0. Grant each processor role bedrock:InvokeModel + bedrock:InvokeModelWithResponseStream on BOTH the inference-profile ARN AND the per-region foundation-model ARNs for us-east-1/us-east-2/us-west-2 (empty-account) -- the us.* profile routes cross-region, so a profile-only grant AccessDenies at runtime. Remove both anthropic-api-key Secret constructs, their grant_read, and the ANTHROPIC_API_KEY_SECRET_ARN env; add BEDROCK_MODEL_ID. The secrets had RemovalPolicy.RETAIN so they are orphaned, not deleted -- flagged in the README for manual post-deploy deletion and key revocation. Add the workorder-email-processor-template-fallback-rate alarm: a FILL(0) + >=10-sample volume-floor MathExpression over the EMF ParseOutcome metric (15-min periods) that pages when the AI-fallback share exceeds 15% sustained, catching Hexagon template drift. ALARM-only SnsAction to site-alerts, no OK action, NOT_BREACHING, matching the existing stack idiom. * Add offline WO parser test suite Cover the deterministic parser with golden-file tests over 55 real scrubbed .eml fixtures (both comment sub-shapes, username Submitted-By, address present/absent, br+CRLF assign addresses), fail-closed validation-gate rules, adversarial and prompt-injection cases that must route to ai_fallback or parse without corrupting other fields, the issue #23 comment_id idempotency invariants, and the Bedrock-fallback dispatch plus EMF-metric emission with a mocked invoke_model. Extend pytest.ini testpaths to discover the co-located suite, and update tests/conftest.load_handler to put a handler's own directory on sys.path so the WO handler's new `from template_parser import ...` resolves under the existing shared handler tests. Point test_local.py at the new template-first + Bedrock flow. Refs: #23 * Document Bedrock migration and WO parse flow in README Record the provider switch to the Bedrock inference profile (no Anthropic API key or Secrets Manager secret, with the retired secrets flagged for manual deletion), the WO deterministic-template-first + AI-fallback flow, the new ParseOutcome EMF metric and template-fallback-rate alarm, the issue #23 comment_id format change, the +00:00 aware-UTC timestamp shift, and offline test instructions. Refs: #23 * Fix f-string lint and formatting in backfill scripts Drop the f prefix from two f-strings that carry no placeholders (F541) and apply ruff format, so `ruff check` / `ruff format --check` pass in CI. * Emit ParseMethod-only EMF set so fallback alarm can fire The fallback-rate alarm queries the ParseOutcome series keyed on ParseMethod alone, but the emitter published only the joint (ParseMethod, TemplateId) dimension set. CloudWatch materializes exactly the listed dimension sets and does not auto-aggregate, so the alarm's series never received data: it evaluated a constant 0 and could never page on template-drift coverage collapse. Publish both ["ParseMethod"] and ["ParseMethod","TemplateId"] and update the EMF regression test to assert both sets are present. * Commit WO parser .eml fixtures for executable coverage The parser test suite globbed for input .eml fixtures that the repo's `*.eml` ignore rule kept uncommitted, so every parametrized golden and fail-closed test collected zero cases and CI could not exercise the deterministic parser that handles 100% of WO email volume. Add a fixtures-only negation to .gitignore and commit the 55 scrubbed positive samples (50 update-plaintext, 5 assign-html) plus 14 ai-fallback and 3 adversarial fixtures. The ai-fallback set covers each fail-closed reason code (subject_no_match, single_space_work_order, malformed_site_code, label_bleed, creation_time_unparseable, wo_id_mismatch, missing_required_field) and the adversarial set proves the parser is total and confines prompt-injection payloads to comment_text without steering the structured fields. * feature: Add PO template parser scaffold and design doc Mirror WO PR #99's template-first approach for the Coupa PO processor. Two templates identified from a full 3,448-email triage: - coupa_new_po (95.5%): scaffolded; fails closed to the LLM until extract_new_po lands. - coupa_cancellation (2.9%): implemented. Nested contract with recursive validation, Decimal money, and a fail-closed gate. Derived fields (site_code/trade/fiscal_year) are deferred to a shared post-stage. Comments, revisions, multi-line, and non-USD emails fall back to Bedrock. docs/po-template-parser.md records the investigation, decisions, and remaining work. Signed-off-by: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> * Implement PO new_po extraction and value-level gate Replace the extract_new_po scaffold stub with the full section-windowed extractor (duplicate-label anchoring, sentinel ship-to, label-keyed U+2022 bullet split, Decimal money from three anchored contexts only) and add value-level gate rules V1-V13. Both new_po_not_implemented scaffold guards are removed; rules 6-8 (unrecognized_status, multiline_unsupported, non_usd) go live. The gate re-derives every byte proof from the email body so an extractor bug cannot vouch for itself: amount re-serialization with a digit/comma border check (the thousands-separator truncation kill switch), sum(lines)==total against both Total blocks, anchor/supplier identity proofs, USPS address shape on the raw pre-enrichment zip, bullet label discipline, and sentinel/artifact hygiene. Any failure falls closed to the LLM; a validation failure is never a parsed result. Refs: #99 * Wire template-first parse into PO handler with EMF metric Run try_deterministic_parse ahead of the Bedrock extractor and fall back only on a miss/invalid (fail-closed) result. The shared enrich_parsed post-stage and the save_cancellation/save_revision/ save_new_po routing are untouched, so both paths write identical DynamoDB shapes and the po-ingest-site-extractor stream contract is preserved. Each record emits one ParseMethod EMF line (Seahaven/PoIngest/ ParseOutcome, dimension sets [ParseMethod] and [ParseMethod,TemplateId], ReasonCode/po_number ride-alongs) mirroring the WO idiom. The metric fires before the Bedrock call so a Bedrock-side error still records the ai_fallback outcome. Refs: #99 * Add PO fallback-rate alarm retuned for ~57 emails/day The WO alarm's 15-min period and >=10-sample floor assume ~760/day and would be structurally dead at PO volume (a 15-min period holds ~0.6 emails, so the floor is never met). Retune: 6-hour periods (~14.25 expected emails), IF((fb+tmpl)>=8,...) volume floor so a single email can never breach a datapoint (1/8 = 12.5% < 20%), threshold >20% against a ~1% expected baseline, eval 4 / datapoints 2 (24h span) so noise self-clears while total template drift pages within ~12h. No element-wise MAX in the math expression (post-#102 rule); ALARM-only SnsAction to site-alerts, NOT_BREACHING. Gated with 'npx cdk synth po-ingest'. Also add template_parser.py to the bundling cp list -- without it every deployed invocation would ImportError (unit tests cannot catch an asset-bundling omission). Refs: #99, #102 * Add offline PO parser suite with scrubbed fixture corpus 132 tests: golden-file comparison for all 25 positive fixtures (17 single-line new-PO + 8 cancellations, Decimal-exact via parse_float=Decimal), every fail-closed gate reason code covered (body-level triggers via 17 synthetic adversarial .eml mutations, candidate-level via direct validate() unit tests), real multi-line and comment/non-Coupa fallback fixtures, dual line-ending parse identity, two-path enrich/save parity (site-extractor stream guard), V10 URL-id corpus sweep, fixture hygiene (ses_auth pass + scrub-marker leak sweep), and Bedrock dispatch/EMF assertions. The suite loads handler/template_parser via importlib under unique module names and binds the handler's bare sibling imports around exec (tests/conftest.py load_handler gets the same treatment) -- the WO suite caches bare 'handler'/'template_parser' names in sys.modules, and bare imports here would silently bind to the wrong pipeline. moto is imported before the handler so its botocore stubber hook precedes boto3 session creation (the PO conftest chain now loads at pytest session start). Fixtures are scrubbed real S3 samples: transport/auth header values replaced with same-shape placeholders (structure kept so ses_auth still passes), per-file digit ciphers, amounts remapped with sum==total re-established. The .gitignore exception is scoped to the PO fixtures path only. Refs: #99 * Document PO template-first parser and retuned alarm README: PO flow is now template-first with Bedrock fallback; parser/gate section mirroring the WO writeup; Seahaven/PoIngest ParseOutcome namespace and the fallback-rate alarm numbers with their volume justification (deliberately not WO's settings); test-suite and repo-layout updates. Design doc: mark PR #1 complete in progress/checklist sections; document the six value-level gate reason codes and the scaffold guard removal; correct the stale data-access note (default CLI session is 328440206208) and note the ~90-day S3 lifecycle aging of the corpus; record the 2.3 layout addendum (leading Supplier bullet segment, EA evidence lines, summary unit-price tokens, decode-path line endings), the fixture-build pins (address join convention, quantity/unit/price source), the V10 sweep outcome, and resolutions for open questions Q3/Q6. Cross-family review and the Confluence architecture-map update are flagged outstanding for merge. Refs: #99 * Record cross-family review outcome for handler wiring GPT-4.1 cross_review.py run against the real handler diff returned no BLOCK and no security findings; both FIX items verified as no-change-needed (fallback logging already correct; non-dict AI output is the pre-existing issue #101 pattern this PR deliberately does not touch). Refs: #99 * Pin line-item currency to USD in the PO gate The non_usd rule only checked the Total-block top-level currency, so a new_po whose line item read 'for 55,206.00 CAD' under a USD Total block still template-parsed as ok -- a fail-open hole in the fail-closed gate. Every line item's captured currency and its re-derived body token must now byte-equal the proven-USD top-level currency; covered by a line-level CAD adversarial fixture (the existing adv-non-usd only exercised the Total-block variant) and a candidate-mutation unit test. * Scrub residual transport tokens from PO fixtures The first-pass harvest scrub sanitized only the primary SES/DKIM header blocks, leaving the real SES Feedback-ID sender-identity hash in 49 committed fixtures and, on the two non-Coupa fixtures, an embedded second SES block's X-Ses-Receipt, the Exchange cross-tenant UPN ciphertext, and Gmail ARC fh= / X-Gm-* tokens -- exactly the token classes the PR #99 fixture lesson requires placeholdered. Replace each with a same-shape ScrubbedFixture value (byte-safe, CRLF and folding preserved) so header structure and ses_auth behavior are unchanged. * Converge quantity/price to Decimal on both paths EXTRACTION_PROMPT declares quantity and price as JSON strings, so a prompt-obedient Bedrock response stores DynamoDB Strings where the template parser stores Numbers -- divergent attribute types for the same email on the purchase-orders stream. Coerce numeric strings to Decimal in the shared enrich_parsed post-stage (thousands-separator safe; non-numeric strings kept verbatim) so both paths converge; prompt rewording itself remains PR #2 scope. The two-path parity test was circular -- it replayed the parser- derived golden as 'the LLM output', so it could never see the type divergence. It now feeds a prompt-shaped payload (string quantity/ price, LLM-filled site_code) through enrich_parsed and save_new_po, and the fixture-hygiene test now asserts the scrubbed transport-token header classes so fixture regressions are caught. * Coerce bare-int quantity/price to Decimal in enrich_parsed GPT-4.1 cross-family review of the final PR diff (no BLOCK) flagged residual type drift: parse_float=Decimal rules out floats on the LLM path, but a bare JSON int survived as Python int. Coerce it so both parse paths emit one canonical Decimal type. * Scrub fixture-body PII and harden cancellation gate (sec review) /sh-security-review of PR #105 (5 fresh-context detectors + proof-or-kill verifier) confirmed two diff-introduced findings; both fixed here. F3 (medium, real PII in new fixtures): the harvest scrub replaced header tokens but left real third-party PII in message BODIES -- an Amazon contact's name/phone/personal email in non-coupa-02.eml and an internal t.corp.amazon.com ticket URL in comment-02.eml, plus real submitter/attn names recurring across the new_po corpus. Replaced every personal name, phone, personal email, and internal URL with synthetic placeholders (QP-soft-wrap aware) across both .eml bodies and expected goldens. Extended test_fixture_hygiene to scan BODIES (phone shapes, corp URLs, the leaked tokens), closing the header-only gap that let this through. F1 (medium, cancellation gate): _CANCELLATION_SUBJECT was unanchored and matched with .search(), unlike the anchored new_po pattern -- a subject merely ending with the cancellation phrase could be routed to the sticky- Cancelled write. Fully anchored it and switched to .match, and added a body-corroboration gate (the real Coupa body independently restates 'Purchase Order #<po> ... has been cancelled'); a near-miss/misrouted subject whose body does not corroborate now fails closed to the LLM (new reason code cancellation_body_unconfirmed). Pre-existing (advisory, not this PR): the LLM-fallback else->save_new_po dispatch and undelimited extraction prompt (issue #101 family) are byte-identical to main and unchanged here. 401 tests pass; ruff/format clean; cdk synth po-ingest clean. --------- Signed-off-by: Adam Moussa <166072409+amoussa1229@users.noreply.github.com>
583 lines
22 KiB
Python
583 lines
22 KiB
Python
"""
|
||
PO email processor Lambda.
|
||
|
||
Triggered by S3 events when SES delivers a Coupa PO email.
|
||
Parses the raw email, tries the deterministic template parser first, falls
|
||
back to Claude on Bedrock for structured extraction on a miss/invalid result,
|
||
then writes the result to the purchase-orders DynamoDB table.
|
||
"""
|
||
|
||
import email
|
||
import json
|
||
import logging
|
||
import os
|
||
import re
|
||
from datetime import datetime, timezone
|
||
from decimal import Decimal, InvalidOperation
|
||
from email import policy
|
||
|
||
import boto3
|
||
from ses_auth import authenticate_inbound_email
|
||
from template_parser import try_deterministic_parse
|
||
|
||
logger = logging.getLogger()
|
||
logger.setLevel(logging.INFO)
|
||
|
||
s3 = boto3.client("s3")
|
||
dynamodb = boto3.resource("dynamodb")
|
||
bedrock = boto3.client("bedrock-runtime")
|
||
|
||
PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders")
|
||
BEDROCK_MODEL_ID = os.environ.get(
|
||
"BEDROCK_MODEL_ID", "us.anthropic.claude-haiku-4-5-20251001-v1:0"
|
||
)
|
||
|
||
# CloudWatch EMF namespace/metric for the parse-outcome metric (the PO
|
||
# fallback-rate alarm in cdk/po_stack.py reads the ["ParseMethod"] series).
|
||
METRIC_NAMESPACE = "Seahaven/PoIngest"
|
||
METRIC_NAME = "ParseOutcome"
|
||
|
||
# "Cancelled" is a sticky, authoritative status: once a PO reaches it, a later
|
||
# new_po/revision may enrich other fields but must never move it back to a
|
||
# non-cancelled status.
|
||
CANCELLED_STATUS = "Cancelled"
|
||
|
||
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 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,
|
||
}
|
||
|
||
|
||
def extract_with_claude(email_data: dict) -> dict:
|
||
"""Send parsed email to Claude on Bedrock for structured extraction."""
|
||
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']}"
|
||
)
|
||
|
||
resp = bedrock.invoke_model(
|
||
modelId=BEDROCK_MODEL_ID,
|
||
body=json.dumps(
|
||
{
|
||
"anthropic_version": "bedrock-2023-05-31",
|
||
"max_tokens": 2048,
|
||
"messages": [
|
||
{
|
||
"role": "user",
|
||
"content": f"{EXTRACTION_PROMPT}\n\nEMAIL:\n{email_text}",
|
||
}
|
||
],
|
||
}
|
||
),
|
||
)
|
||
response_text = json.loads(resp["body"].read())["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)
|
||
|
||
# parse_float=Decimal is CRITICAL: DynamoDB rejects Python floats.
|
||
return json.loads(response_text.strip(), parse_float=Decimal)
|
||
|
||
|
||
def _emit_parse_method_metric(method, template_id, reason_code, po_number):
|
||
"""Emit one CloudWatch EMF line recording the parse outcome.
|
||
|
||
Zero-latency (no PutMetricData API call): the extraction path is async and
|
||
the role already has logs:PutLogEvents. ParseMethod/TemplateId are the only
|
||
promoted (dimensioned) fields to keep cardinality low; ReasonCode and
|
||
po_number ride along as Logs-Insights-queryable properties.
|
||
|
||
Two dimension sets are published: ["ParseMethod"] (aggregated across all
|
||
template ids -- the series the fallback-rate alarm queries) AND
|
||
["ParseMethod", "TemplateId"] (per-template breakdown for Logs Insights /
|
||
dashboards). CloudWatch materializes only the exact dimension sets listed
|
||
here and does NOT auto-aggregate, so the alarm's single-dimension query
|
||
would receive no data unless ["ParseMethod"] is emitted explicitly."""
|
||
emf = {
|
||
"_aws": {
|
||
# EMF requires Timestamp (epoch ms); without it CloudWatch may not
|
||
# extract the metric datapoint from the log event.
|
||
"Timestamp": int(datetime.now(timezone.utc).timestamp() * 1000),
|
||
"CloudWatchMetrics": [
|
||
{
|
||
"Namespace": METRIC_NAMESPACE,
|
||
"Dimensions": [["ParseMethod"], ["ParseMethod", "TemplateId"]],
|
||
"Metrics": [{"Name": METRIC_NAME, "Unit": "Count"}],
|
||
}
|
||
],
|
||
},
|
||
"ParseMethod": method,
|
||
"TemplateId": template_id or "unknown",
|
||
"ReasonCode": reason_code or "ok",
|
||
"po_number": po_number or "",
|
||
METRIC_NAME: 1,
|
||
}
|
||
print(json.dumps(emf))
|
||
|
||
|
||
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.now(timezone.utc).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"])
|
||
|
||
# Canonical numeric type for quantity/price across BOTH parse paths: the
|
||
# template parser emits Decimal (DynamoDB Number) while EXTRACTION_PROMPT
|
||
# asks the LLM for these two fields as JSON strings (which the Bedrock
|
||
# json.loads leaves as str -> DynamoDB String). Coercing here -- in the
|
||
# SHARED post-stage -- converges the attribute type to Number for
|
||
# equivalent parsed dicts, preserving the two-path parity contract on the
|
||
# purchase-orders table stream. Prompt rewording itself is PR #2 scope.
|
||
# A non-numeric string is left verbatim (still stored, as a String) --
|
||
# dropping it would lose LLM-extracted evidence.
|
||
for item in parsed.get("line_items") or []:
|
||
if not isinstance(item, dict):
|
||
continue
|
||
for field in ("quantity", "price"):
|
||
value = item.get(field)
|
||
if isinstance(value, str):
|
||
try:
|
||
item[field] = Decimal(value.replace(",", "").strip())
|
||
except InvalidOperation:
|
||
pass
|
||
elif isinstance(value, (int, float)) and not isinstance(value, bool):
|
||
# parse_float=Decimal means floats can't occur on the LLM path,
|
||
# but a bare JSON int would slip through as Python int; coerce
|
||
# so both paths emit one canonical Decimal type (cross-review FIX).
|
||
item[field] = Decimal(str(value))
|
||
return parsed
|
||
|
||
|
||
def _write_fields(po_number: str, fields: dict, *, guard_cancelled: bool):
|
||
"""SET the given non-null fields on a PO record via update_item.
|
||
|
||
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).
|
||
|
||
When ``guard_cancelled`` is True the write carries a ConditionExpression that
|
||
only permits it while the record is not already Cancelled. The condition is
|
||
evaluated atomically by DynamoDB at write time, so a cancellation that lands
|
||
first always wins — there is no read-then-write TOCTOU window. A failed guard
|
||
raises ConditionalCheckFailedException for the caller to handle.
|
||
"""
|
||
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
|
||
|
||
params = {
|
||
"Key": {"po_number": po_number},
|
||
"UpdateExpression": "SET " + ", ".join(set_parts),
|
||
"ExpressionAttributeNames": attr_names,
|
||
"ExpressionAttributeValues": attr_values,
|
||
}
|
||
if guard_cancelled:
|
||
params["ExpressionAttributeValues"][":__cancelled_marker"] = CANCELLED_STATUS
|
||
params["ConditionExpression"] = (
|
||
"attribute_not_exists(po_status) OR po_status <> :__cancelled_marker"
|
||
)
|
||
|
||
table.update_item(**params)
|
||
|
||
|
||
def _merge_update(po_number: str, fields: dict):
|
||
"""Merge (SET-only) the given fields onto a PO record, keeping Cancelled sticky.
|
||
|
||
Only the fields supplied are written; absent fields are left untouched. The
|
||
record is created if it does not exist (DynamoDB update_item upsert).
|
||
|
||
"Cancelled" is a sticky, authoritative status. When the incoming payload
|
||
carries a non-cancelled ``po_status``, the write is guarded by a
|
||
ConditionExpression so the status is only applied while the record is not
|
||
already Cancelled — enforced atomically at write time, eliminating the
|
||
read-then-write TOCTOU where a concurrently-landing cancellation could be
|
||
silently un-cancelled. If the guard fails (the PO is already Cancelled), the
|
||
same fields are re-written WITHOUT po_status/cancelled_at and
|
||
unconditionally, so the other fields still merge while the Cancelled status
|
||
stays intact.
|
||
|
||
A payload with no ``po_status``, or one whose status is already "Cancelled",
|
||
needs no guard — a plain merge is correct. This is what keeps legitimate
|
||
status updates (non-cancelled PO) and status-less revisions from ever being
|
||
dropped: the status is only ever suppressed on a true un-cancel transition.
|
||
"""
|
||
incoming_status = fields.get("po_status")
|
||
if incoming_status is None or incoming_status == CANCELLED_STATUS:
|
||
_write_fields(po_number, fields, guard_cancelled=False)
|
||
return
|
||
|
||
try:
|
||
_write_fields(po_number, fields, guard_cancelled=True)
|
||
except dynamodb.meta.client.exceptions.ConditionalCheckFailedException:
|
||
logger.info(
|
||
f"PO {po_number} is Cancelled; suppressing incoming "
|
||
f"po_status={incoming_status!r} and merging remaining fields"
|
||
)
|
||
enrich_fields = {
|
||
k: v for k, v in fields.items() if k not in ("po_status", "cancelled_at")
|
||
}
|
||
_write_fields(po_number, enrich_fields, guard_cancelled=False)
|
||
|
||
|
||
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. "Cancelled" is a sticky
|
||
status enforced atomically inside _merge_update: if the PO was already
|
||
cancelled, the new_po backfills its remaining fields (supplier, line_items,
|
||
amounts) but never un-cancels it.
|
||
"""
|
||
po_number = parsed["po_number"]
|
||
fields = {k: v for k, v in parsed.items() if v is not None}
|
||
_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.
|
||
|
||
"Cancelled" is a sticky status: a revision may enrich a cancelled PO's fields
|
||
but must never move it to a non-cancelled status. That invariant is enforced
|
||
atomically inside _merge_update and applies ONLY to the un-cancel transition —
|
||
a revision that carries no status change, or one targeting a non-cancelled PO,
|
||
updates po_status normally.
|
||
"""
|
||
po_number = parsed["po_number"]
|
||
fields = {k: v for k, v in parsed.items() if v is not 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.now(timezone.utc).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()
|
||
|
||
# Fail-closed sender authentication (INFRA-107): only mail with an
|
||
# SES-stamped dkim=pass verdict for an allowlisted domain may create
|
||
# or update purchase orders. Rejected mail is logged and skipped
|
||
# without erroring the invocation (no retries / DLQ spam).
|
||
if not authenticate_inbound_email(raw_email, s3_key):
|
||
continue
|
||
|
||
email_data = parse_raw_email(raw_email)
|
||
logger.info(f"Subject: {email_data['subject']}")
|
||
|
||
# Deterministic template parse first; fall back to the Bedrock AI
|
||
# extractor only on a miss or an invalid (fail-closed) result. The
|
||
# metric is emitted before the Bedrock call, so a Bedrock-side error
|
||
# still records the ai_fallback outcome (the invocation then errors
|
||
# into the errors alarm / DLQ as before).
|
||
parsed, parse_method, template_id, parse_reason = try_deterministic_parse(
|
||
email_data
|
||
)
|
||
_emit_parse_method_metric(
|
||
parse_method,
|
||
template_id,
|
||
parse_reason,
|
||
parsed.get("po_number") if parsed else None,
|
||
)
|
||
if parsed is None:
|
||
parsed = extract_with_claude(email_data)
|
||
logger.info(
|
||
f"Parsed ({parse_method}/{template_id}/{parse_reason}): "
|
||
f"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"}
|