Merge pull request #2 from Sea-Haven-Industries/feature/verified-sites-pipeline

Add verified-sites pipeline with address fallback
This commit is contained in:
Adam Moussa 2026-04-30 15:06:04 -04:00 • committed by GitHub
commit 2b8e121413
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
5 changed files with 500 additions and 8 deletions

View file

@ -10,15 +10,36 @@ Coupa purchase-order email ingestion pipeline. SES receives Amazon PO emails, Cl
4. The Lambda parses the email, sends it to Claude Haiku 4.5 for structured JSON extraction, and writes to DynamoDB. 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: new_po` — conditional `PutItem` on `purchase-orders` (idempotent on `po_number`).
- `email_type: cancellation` — `UpdateItem` marking the existing row `Cancelled`. - `email_type: cancellation` — `UpdateItem` marking the existing row `Cancelled`.
5. LedgerFlow consumes the table via DynamoDB Streams → `po-sync`. This repo only writes. 5. DynamoDB Streams (NEW_AND_OLD_IMAGES) 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).
A separate `po-web-ui` Lambda (Function URL, unauthenticated) renders a simple HTML dashboard scanning the table. A separate `po-web-ui` Lambda (Function URL, unauthenticated) renders a simple HTML dashboard scanning the table.
### Verified-sites pipeline
The `po-ingest-site-extractor` Lambda is triggered by the DynamoDB Stream on every PO INSERT/MODIFY. It:
1. Extracts an Amazon facility site code from `ship_to.name` using a regex cascade (parentheses, `LLC - CODE`, `Station CODE`, `DS - CODE`) with a fallback to the first `line_items` description.
2. Parses `ship_to.address` into structured fields (street, city, state, zip).
3. Upserts to the `verified-sites` DynamoDB table — atomically increments `poCount` and appends the PO number to `sourcePOs`.
POs with no extractable site code fall through to an address reverse-lookup against the verified-sites cache (normalized street + zip). If still unresolved, the PO is written to the `pending-site-review` table for manual verification against Payee Central.
Backfill stats (initial run): 14,825 POs scanned → 9,900 with extractable site codes → 1,100 unique sites.
## Architecture ## Architecture
- **IaC:** AWS CDK (Python), stack name `PoIngestStack`, region `us-east-1`. - **IaC:** AWS CDK (Python), stack name `PoIngestStack`, region `us-east-1`.
- **Lambdas:** `po-email-processor` (S3-triggered) and `po-web-ui` (Function URL). Python 3.12, 256 MB, 60s timeout. - **Lambdas** (all Python 3.12, arm64, 60-day log retention):
- **Storage:** S3 `po-ingest-emails-{AccountId}` with 90-day lifecycle expiry; DynamoDB `purchase-orders` (shared, not owned by this stack). - `po-email-processor` — S3-triggered, parses PO emails via Claude Haiku.
- `po-web-ui` — Function URL, HTML dashboard.
- `po-ingest-site-extractor` — DynamoDB Streams-triggered, extracts site addresses.
- **Storage:**
- S3 `po-ingest-emails-{AccountId}` — 90-day lifecycle expiry.
- DynamoDB `purchase-orders` — owned by this stack, Streams enabled (NEW_AND_OLD_IMAGES).
- DynamoDB `verified-sites` — PK `siteCode`, GSI `by-state` on `state`.
- DynamoDB `pending-site-review` — PK `po_number`. POs with no extractable site code and no address match, awaiting manual Payee Central verification.
- **Secrets:** Anthropic API key in Secrets Manager at `po-ingest/anthropic-api-key`. - **Secrets:** Anthropic API key in Secrets Manager at `po-ingest/anthropic-api-key`.
- **SES:** adds the `PoEmailRule` to the existing `INBOUND_MAIL` receipt rule set (shared with `workorder-ingest`). - **SES:** adds the `PoEmailRule` to the existing `INBOUND_MAIL` receipt rule set (shared with `workorder-ingest`).
@ -53,3 +74,13 @@ python scripts/reprocess.py --execute # invokes po-email-processor for each
``` ```
Inserts are conditional on `po_number`, so re-processing existing POs is a no-op. Inserts are conditional on `po_number`, so re-processing existing POs is a no-op.
## Backfilling verified sites
The stream Lambda handles all future POs automatically. To backfill from historical PO data (one-time):
```bash
python scripts/backfill_sites.py
```
Uses the same extraction logic as the Lambda. Idempotent — safe to re-run.

View file

@ -7,6 +7,8 @@ from aws_cdk import (
Stack, Stack,
aws_dynamodb as dynamodb, aws_dynamodb as dynamodb,
aws_lambda as lambda_, aws_lambda as lambda_,
aws_lambda_event_sources as lambda_event_sources,
aws_logs as logs,
aws_s3 as s3, aws_s3 as s3,
aws_s3_notifications as s3n, aws_s3_notifications as s3n,
aws_ses as ses, aws_ses as ses,
@ -31,11 +33,19 @@ class PoIngestStack(Stack):
], ],
) )
# --- Reference existing purchase-orders DynamoDB table --- # --- Purchase-orders DynamoDB table ---
# This table is shared with LedgerFlow (DynamoDB Streams → po-sync). # Owned by this stack. Streams enabled for the site-extractor pipeline.
# We write to it; LedgerFlow reads from it. # Other stacks (seahaven-slack-bot) reference this table via fromTableName().
po_table = dynamodb.Table.from_table_name( po_table = dynamodb.Table(
self, "PurchaseOrdersTable", "purchase-orders", self, "PurchaseOrdersTable",
table_name="purchase-orders",
partition_key=dynamodb.Attribute(
name="po_number",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
stream=dynamodb.StreamViewType.NEW_AND_OLD_IMAGES,
) )
# --- Secrets Manager for Anthropic API key --- # --- Secrets Manager for Anthropic API key ---
@ -50,10 +60,12 @@ class PoIngestStack(Stack):
self, "EmailProcessor", self, "EmailProcessor",
function_name="po-email-processor", function_name="po-email-processor",
runtime=lambda_.Runtime.PYTHON_3_12, runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler", handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/email_processor/package"), code=lambda_.Code.from_asset("../lambdas/email_processor/package"),
timeout=Duration.seconds(60), timeout=Duration.seconds(60),
memory_size=256, memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={ environment={
"PO_TABLE": "purchase-orders", "PO_TABLE": "purchase-orders",
"ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn, "ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn,
@ -94,10 +106,12 @@ class PoIngestStack(Stack):
self, "WebUI", self, "WebUI",
function_name="po-web-ui", function_name="po-web-ui",
runtime=lambda_.Runtime.PYTHON_3_12, runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler", handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/web_ui"), code=lambda_.Code.from_asset("../lambdas/web_ui"),
timeout=Duration.seconds(60), timeout=Duration.seconds(60),
memory_size=256, memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={ environment={
"PO_TABLE": "purchase-orders", "PO_TABLE": "purchase-orders",
}, },
@ -111,3 +125,73 @@ class PoIngestStack(Stack):
) )
cdk.CfnOutput(self, "WebUIUrl", value=web_url.url, description="PO Dashboard URL") cdk.CfnOutput(self, "WebUIUrl", value=web_url.url, description="PO Dashboard URL")
# --- Verified sites table (extracted from PO ship-to addresses) ---
verified_sites_table = dynamodb.Table(
self, "VerifiedSitesTable",
table_name="verified-sites",
partition_key=dynamodb.Attribute(
name="siteCode",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
verified_sites_table.add_global_secondary_index(
index_name="by-state",
partition_key=dynamodb.Attribute(
name="state",
type=dynamodb.AttributeType.STRING,
),
projection_type=dynamodb.ProjectionType.ALL,
)
# --- Site extractor Lambda (DynamoDB Streams → verified-sites) ---
site_extractor = lambda_.Function(
self, "SiteExtractor",
function_name="po-ingest-site-extractor",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/site_extractor"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={
"VERIFIED_SITES_TABLE": verified_sites_table.table_name,
"PENDING_REVIEW_TABLE": "pending-site-review",
},
)
verified_sites_table.grant_read_write_data(site_extractor)
site_extractor.add_event_source(
lambda_event_sources.DynamoEventSource(
po_table,
starting_position=lambda_.StartingPosition.TRIM_HORIZON,
batch_size=10,
max_batching_window=Duration.seconds(30),
bisect_batch_on_error=True,
retry_attempts=3,
)
)
cdk.CfnOutput(self, "VerifiedSitesTableName",
value=verified_sites_table.table_name,
description="Verified site addresses extracted from POs",
)
# --- Pending site review table (POs with no extractable site code) ---
pending_review_table = dynamodb.Table(
self, "PendingSiteReviewTable",
table_name="pending-site-review",
partition_key=dynamodb.Attribute(
name="po_number",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
pending_review_table.grant_read_write_data(site_extractor)
verified_sites_table.grant_read_data(site_extractor)

View file

@ -52,9 +52,14 @@ Analyze the following email and extract structured data. Return ONLY valid JSON
"supplier": { "supplier": {
"name": "string or null" "name": "string or null"
}, },
"site_code": "string or null",
"ship_to": { "ship_to": {
"name": "string or null", "name": "string or null",
"address": "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", "location_code": "string or null",
"attn": "string or null" "attn": "string or null"
}, },
@ -80,6 +85,13 @@ Rules:
- Extract the PO number from the email (e.g., "2D-18206023") - Extract the PO number from the email (e.g., "2D-18206023")
- Extract all line items with their descriptions, amounts, and metadata - 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 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
- total_amount should be the numeric total in USD - total_amount should be the numeric total in USD
- If a field is not present in the email, set it to null - 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 - Do NOT invent or infer data that is not explicitly in the email

View file

@ -0,0 +1,289 @@
"""
DynamoDB Streams processor that extracts Amazon site codes and addresses
from purchase order records and upserts them into the verified-sites table.
"""
import logging
import os
import re
from datetime import datetime, timezone
import boto3
from boto3.dynamodb.types import TypeDeserializer
logger = logging.getLogger()
logger.setLevel(logging.INFO)
dynamodb = boto3.client("dynamodb")
deserializer = TypeDeserializer()
VERIFIED_SITES_TABLE = os.environ.get("VERIFIED_SITES_TABLE", "verified-sites")
PENDING_REVIEW_TABLE = os.environ.get("PENDING_REVIEW_TABLE", "pending-site-review")
SITE_CODE_PATTERN = re.compile(r"[A-Z]{2,4}\d{1,2}")
# Cold-start cache: normalized (street, zip) → siteCode
_address_cache: dict[tuple[str, str], str] = {}
_cache_loaded = False
def deserialize_image(image: dict) -> dict:
return {k: deserializer.deserialize(v) for k, v in image.items()}
def extract_site_code(record: dict) -> str | None:
# Prefer the explicit site_code field (set by updated extraction prompt)
direct = (record.get("site_code") or "").strip().upper()
if direct and SITE_CODE_PATTERN.match(direct):
return direct
ship_to = record.get("ship_to") or {}
ship_to_name = ship_to.get("name", "")
if ship_to_name:
m = re.search(r"\(([A-Z0-9]{3,5})\)", ship_to_name)
if m and SITE_CODE_PATTERN.match(m.group(1)):
return m.group(1)
m = re.search(r"(?:LLC|Inc)\s*-\s*([A-Z0-9]{3,5})\b", ship_to_name)
if m and SITE_CODE_PATTERN.match(m.group(1)):
return m.group(1)
m = re.search(
r"(?:Station|DS)\s*[-–]?\s*([A-Z0-9]{3,5})\b", ship_to_name
)
if m and SITE_CODE_PATTERN.match(m.group(1)):
return m.group(1)
m = re.match(r"^([A-Z0-9]{3,5})\s*-\s*Amazon", ship_to_name)
if m and SITE_CODE_PATTERN.match(m.group(1)):
return m.group(1)
for item in record.get("line_items") or []:
desc = item.get("description", "")
m = re.match(r"^([A-Z]{1,4}[0-9]{1,2}|[A-Z]{4,5})\b", desc)
if m and SITE_CODE_PATTERN.match(m.group(1)):
return m.group(1)
return None
def parse_address(record: dict) -> dict:
result = {
"address": None,
"city": None,
"state": None,
"zip": None,
"fullAddress": None,
}
ship_to = record.get("ship_to") or {}
# Prefer structured fields if the extraction prompt provided them
if ship_to.get("street") and ship_to.get("state"):
result["address"] = ship_to["street"]
result["city"] = ship_to.get("city")
result["state"] = ship_to["state"]
result["zip"] = ship_to.get("zip")
parts = [p for p in [result["address"], result["city"], result["state"], result["zip"]] if p]
result["fullAddress"] = ", ".join(parts)
return result
address_text = ship_to.get("address", "")
if not address_text:
return result
clean = re.sub(r",?\s*United States\s*$", "", address_text.strip())
if not clean:
return result
# Match "city, ST ZIP" at the end of the string
m = re.search(r",\s*([^,]+?),\s+([A-Z]{2})\s+(\d{5}(?:-\d{4})?)\s*$", clean)
if m:
result["address"] = clean[: m.start()].strip()
result["city"] = m.group(1).strip()
result["state"] = m.group(2)
result["zip"] = m.group(3)
result["fullAddress"] = f"{result['address']}, {result['city']}, {result['state']} {result['zip']}"
return result
# Fallback: try newline-separated format
lines = [l.strip() for l in clean.split("\n") if l.strip()]
city_re = re.compile(r"^(.+?),\s+([A-Z]{2})\s+(\d{5}(?:-\d{4})?)")
for i, line in enumerate(lines):
m = city_re.match(line)
if m:
result["address"] = ", ".join(lines[:i]) if i > 0 else None
result["city"] = m.group(1)
result["state"] = m.group(2)
result["zip"] = m.group(3)
street = result["address"] or result["city"]
result["fullAddress"] = f"{street}, {result['city']}, {result['state']} {result['zip']}"
return result
result["address"] = clean
return result
def upsert_site(
site_code: str,
address: dict,
po_number: str,
location_code: str | None,
):
now = datetime.now(timezone.utc).isoformat()
expr_names = {}
expr_values = {
":now": {"S": now},
":one": {"N": "1"},
":po_set": {"SS": [po_number]},
}
set_parts = ["lastSeenAt = :now"]
if address.get("address"):
set_parts.append("#addr = :addr")
expr_names["#addr"] = "address"
expr_values[":addr"] = {"S": address["address"]}
if address.get("city"):
set_parts.append("city = :city")
expr_values[":city"] = {"S": address["city"]}
if address.get("state"):
set_parts.append("#st = :state_val")
expr_names["#st"] = "state"
expr_values[":state_val"] = {"S": address["state"]}
if address.get("zip"):
set_parts.append("zip = :zip")
expr_values[":zip"] = {"S": address["zip"]}
if address.get("fullAddress"):
set_parts.append("fullAddress = :full")
expr_values[":full"] = {"S": address["fullAddress"]}
if location_code:
set_parts.append("locationCode = :loc")
expr_values[":loc"] = {"S": location_code}
update_expr = f"SET {', '.join(set_parts)} ADD poCount :one, sourcePOs :po_set"
kwargs = {
"TableName": VERIFIED_SITES_TABLE,
"Key": {"siteCode": {"S": site_code}},
"UpdateExpression": update_expr,
"ExpressionAttributeValues": expr_values,
}
if expr_names:
kwargs["ExpressionAttributeNames"] = expr_names
dynamodb.update_item(**kwargs)
def _normalize_street(street: str) -> str:
s = street.strip().upper()
s = re.sub(r"[.,#]", "", s)
s = re.sub(r"\s+", " ", s)
return s
def load_address_cache():
global _address_cache, _cache_loaded
if _cache_loaded:
return
paginator = dynamodb.get_paginator("scan")
for page in paginator.paginate(
TableName=VERIFIED_SITES_TABLE,
ProjectionExpression="siteCode, #addr, zip",
ExpressionAttributeNames={"#addr": "address"},
):
for item in page.get("Items", []):
site_code = item.get("siteCode", {}).get("S")
street = item.get("address", {}).get("S")
zip_code = item.get("zip", {}).get("S")
if site_code and street and zip_code:
key = (_normalize_street(street), zip_code)
_address_cache[key] = site_code
_cache_loaded = True
logger.info("Address cache loaded: %d entries", len(_address_cache))
def lookup_by_address(address: dict) -> str | None:
street = address.get("address")
zip_code = address.get("zip")
if not street or not zip_code:
return None
return _address_cache.get((_normalize_street(street), zip_code))
def write_pending_review(po_number: str, address: dict, ship_to_name: str | None):
now = datetime.now(timezone.utc).isoformat()
item = {
"po_number": {"S": po_number},
"createdAt": {"S": now},
"status": {"S": "pending"},
}
if ship_to_name:
item["shipToName"] = {"S": ship_to_name}
if address.get("fullAddress"):
item["fullAddress"] = {"S": address["fullAddress"]}
if address.get("address"):
item["address"] = {"S": address["address"]}
if address.get("city"):
item["city"] = {"S": address["city"]}
if address.get("state"):
item["state"] = {"S": address["state"]}
if address.get("zip"):
item["zip"] = {"S": address["zip"]}
dynamodb.put_item(TableName=PENDING_REVIEW_TABLE, Item=item)
def handler(event, context):
load_address_cache()
processed = 0
extracted = 0
matched_by_address = 0
pending = 0
for record in event.get("Records", []):
if record["eventName"] not in ("INSERT", "MODIFY"):
continue
new_image = record.get("dynamodb", {}).get("NewImage")
if not new_image:
continue
processed += 1
po = deserialize_image(new_image)
po_number = po.get("po_number", "unknown")
address = parse_address(po)
location_code = (po.get("ship_to") or {}).get("location_code")
site_code = extract_site_code(po)
if not site_code:
site_code = lookup_by_address(address)
if site_code:
matched_by_address += 1
logger.info(
"Matched PO %s to site %s via address lookup", po_number, site_code
)
if site_code:
upsert_site(site_code, address, po_number, location_code)
extracted += 1
logger.info("Upserted site %s from PO %s", site_code, po_number)
else:
ship_to_name = (po.get("ship_to") or {}).get("name")
write_pending_review(po_number, address, ship_to_name)
pending += 1
logger.info("PO %s added to pending review (no site code or address match)", po_number)
logger.info(
"Batch complete: %d processed, %d extracted (%d via address), %d pending review",
processed,
extracted,
matched_by_address,
pending,
)

76
scripts/backfill_sites.py Normal file
View file

@ -0,0 +1,76 @@
"""
One-time backfill script: scans purchase-orders and populates verified-sites.
Usage:
python scripts/backfill_sites.py
"""
import sys
import os
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "lambdas", "site_extractor"))
import boto3
from handler import extract_site_code, parse_address, upsert_site
PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders")
REGION = os.environ.get("AWS_DEFAULT_REGION", "us-east-1")
dynamodb_resource = boto3.resource("dynamodb", region_name=REGION)
def scan_all_pos():
table = dynamodb_resource.Table(PO_TABLE)
records = []
last_key = None
while True:
kwargs = {}
if last_key:
kwargs["ExclusiveStartKey"] = last_key
response = table.scan(**kwargs)
records.extend(response.get("Items", []))
last_key = response.get("LastEvaluatedKey")
if not last_key:
break
print(f" Scanned {len(records)} POs so far...")
return records
def main():
print(f"Scanning {PO_TABLE} table...")
pos = scan_all_pos()
print(f"Found {len(pos)} purchase orders")
extracted = 0
skipped = 0
sites_seen = set()
for po in pos:
po_number = po.get("po_number", "unknown")
site_code = extract_site_code(po)
if not site_code:
skipped += 1
continue
address = parse_address(po)
location_code = (po.get("ship_to") or {}).get("location_code")
upsert_site(site_code, address, po_number, location_code)
extracted += 1
sites_seen.add(site_code)
if extracted % 100 == 0:
print(f" Processed {extracted} POs with site codes...")
print(f"\nBackfill complete:")
print(f" Total POs scanned: {len(pos)}")
print(f" POs with site code: {extracted}")
print(f" POs without site code: {skipped}")
print(f" Unique sites upserted: {len(sites_seen)}")
if __name__ == "__main__":
main()