procurement-ingest/lambdas/site_extractor/handler.py
Adam Moussa abdf2aa035
Add CI workflow (#18)
* Add CI workflow and apply ruff formatting

* Disable cdk synth — email_processor uses pre-built package dir

The email_processor Lambda bundles deps into a gitignored package/
directory. cdk synth fails in CI without a build step to recreate it.
Disabling until packaging is standardized.

* Use CDK BundlingOptions for email_processor Lambda packaging

Replaces the pre-built gitignored package/ directory with CDK's
built-in bundling. Deps are now installed inside a Docker container
during cdk synth, so the build works identically locally and in CI.
Re-enables run-cdk-synth in the CI workflow.
2026-05-08 16:01:21 -04:00

298 lines
9.3 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.

"""
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 = [line.strip() for line in clean.split("\n") if line.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,
)