""" 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, )