From ec416079f55fcee6038efe028355660dd3f52048 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 14:26:53 -0400 Subject: [PATCH 1/5] Add verified-sites pipeline via DynamoDB Streams Enable DynamoDB Streams on purchase-orders table and add a site-extractor Lambda that extracts Amazon facility codes and addresses from PO ship-to data, upserting them into a new verified-sites table. Includes a backfill script for existing POs and upgrades existing Lambdas to arm64 + 60-day log retention. --- cdk/stack.py | 78 +++++++++++- lambdas/site_extractor/handler.py | 193 ++++++++++++++++++++++++++++++ scripts/backfill_sites.py | 76 ++++++++++++ 3 files changed, 342 insertions(+), 5 deletions(-) create mode 100644 lambdas/site_extractor/handler.py create mode 100644 scripts/backfill_sites.py diff --git a/cdk/stack.py b/cdk/stack.py index 5e343aa..886dd12 100644 --- a/cdk/stack.py +++ b/cdk/stack.py @@ -7,6 +7,8 @@ from aws_cdk import ( Stack, aws_dynamodb as dynamodb, aws_lambda as lambda_, + aws_lambda_event_sources as lambda_event_sources, + aws_logs as logs, aws_s3 as s3, aws_s3_notifications as s3n, aws_ses as ses, @@ -31,11 +33,19 @@ class PoIngestStack(Stack): ], ) - # --- Reference existing purchase-orders DynamoDB table --- - # This table is shared with LedgerFlow (DynamoDB Streams → po-sync). - # We write to it; LedgerFlow reads from it. - po_table = dynamodb.Table.from_table_name( - self, "PurchaseOrdersTable", "purchase-orders", + # --- Purchase-orders DynamoDB table --- + # Owned by this stack. Streams enabled for the site-extractor pipeline. + # Other stacks (seahaven-slack-bot) reference this table via fromTableName(). + po_table = dynamodb.Table( + 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 --- @@ -50,10 +60,12 @@ class PoIngestStack(Stack): self, "EmailProcessor", function_name="po-email-processor", runtime=lambda_.Runtime.PYTHON_3_12, + architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset("../lambdas/email_processor/package"), timeout=Duration.seconds(60), memory_size=256, + log_retention=logs.RetentionDays.TWO_MONTHS, environment={ "PO_TABLE": "purchase-orders", "ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn, @@ -94,10 +106,12 @@ class PoIngestStack(Stack): self, "WebUI", function_name="po-web-ui", runtime=lambda_.Runtime.PYTHON_3_12, + architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset("../lambdas/web_ui"), timeout=Duration.seconds(60), memory_size=256, + log_retention=logs.RetentionDays.TWO_MONTHS, environment={ "PO_TABLE": "purchase-orders", }, @@ -111,3 +125,57 @@ class PoIngestStack(Stack): ) 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, + }, + ) + + 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", + ) diff --git a/lambdas/site_extractor/handler.py b/lambdas/site_extractor/handler.py new file mode 100644 index 0000000..7d8027f --- /dev/null +++ b/lambdas/site_extractor/handler.py @@ -0,0 +1,193 @@ +""" +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") + +SITE_CODE_PATTERN = re.compile(r"[A-Z]{2,4}\d{1,2}") + + +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: + 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 {} + 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 handler(event, context): + processed = 0 + extracted = 0 + skipped = 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") + + site_code = extract_site_code(po) + if not site_code: + skipped += 1 + logger.info("No site code for PO %s", po_number) + 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 + logger.info("Upserted site %s from PO %s", site_code, po_number) + + logger.info( + "Batch complete: %d processed, %d extracted, %d skipped", + processed, + extracted, + skipped, + ) diff --git a/scripts/backfill_sites.py b/scripts/backfill_sites.py new file mode 100644 index 0000000..121d5ca --- /dev/null +++ b/scripts/backfill_sites.py @@ -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() From 194b1aec44aab75903ce4024f3d09a46a0a29d13 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 14:30:54 -0400 Subject: [PATCH 2/5] Update README with verified-sites pipeline and architecture changes --- README.md | 36 +++++++++++++++++++++++++++++++++--- 1 file changed, 33 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index 44b0bce..e11c80d 100644 --- a/README.md +++ b/README.md @@ -10,15 +10,35 @@ 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. - `email_type: new_po` — conditional `PutItem` on `purchase-orders` (idempotent on `po_number`). - `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. +### 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 are logged and skipped. + +Backfill stats (initial run): 14,825 POs scanned → 9,900 with extractable site codes → 1,100 unique sites. + ## Architecture - **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. -- **Storage:** S3 `po-ingest-emails-{AccountId}` with 90-day lifecycle expiry; DynamoDB `purchase-orders` (shared, not owned by this stack). +- **Lambdas** (all Python 3.12, arm64, 60-day log retention): + - `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`. - **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`). @@ -53,3 +73,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. + +## 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. From 5a86b3c96b3c0d9732ba7a60aea5a3db840cf458 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 14:42:09 -0400 Subject: [PATCH 3/5] Add site_code and structured address fields to extraction prompt The Claude extraction prompt now explicitly asks for site_code (the Amazon facility code) and structured ship_to address fields (street, city, state, zip). The site-extractor Lambda prefers these direct fields when available, falling back to regex for older PO records. --- lambdas/email_processor/handler.py | 12 ++++++++++++ lambdas/site_extractor/handler.py | 16 ++++++++++++++++ 2 files changed, 28 insertions(+) diff --git a/lambdas/email_processor/handler.py b/lambdas/email_processor/handler.py index 4f28086..eea5a8e 100644 --- a/lambdas/email_processor/handler.py +++ b/lambdas/email_processor/handler.py @@ -52,9 +52,14 @@ Analyze the following email and extract structured data. Return ONLY valid JSON "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" }, @@ -80,6 +85,13 @@ Rules: - 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 - 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 diff --git a/lambdas/site_extractor/handler.py b/lambdas/site_extractor/handler.py index 7d8027f..2a187d1 100644 --- a/lambdas/site_extractor/handler.py +++ b/lambdas/site_extractor/handler.py @@ -27,6 +27,11 @@ def deserialize_image(image: dict) -> dict: 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", "") @@ -68,6 +73,17 @@ def parse_address(record: dict) -> dict: } 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 From f5d1eaeb5b94be150bdfaf2e1a57349cb86cd2f4 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 15:01:09 -0400 Subject: [PATCH 4/5] Add address reverse-lookup fallback and pending-site-review table When no site code is found via Claude extraction or regex cascade, the Lambda now checks the PO address against a cold-start cache of verified-sites (normalized street + zip match). If still no match, the PO is written to a new pending-site-review table for manual verification against Payee Central. --- cdk/stack.py | 16 +++++ lambdas/site_extractor/handler.py | 106 ++++++++++++++++++++++++++---- 2 files changed, 109 insertions(+), 13 deletions(-) diff --git a/cdk/stack.py b/cdk/stack.py index 886dd12..cdb0908 100644 --- a/cdk/stack.py +++ b/cdk/stack.py @@ -159,6 +159,7 @@ class PoIngestStack(Stack): log_retention=logs.RetentionDays.TWO_MONTHS, environment={ "VERIFIED_SITES_TABLE": verified_sites_table.table_name, + "PENDING_REVIEW_TABLE": "pending-site-review", }, ) @@ -179,3 +180,18 @@ class PoIngestStack(Stack): 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) diff --git a/lambdas/site_extractor/handler.py b/lambdas/site_extractor/handler.py index 2a187d1..f3a9eff 100644 --- a/lambdas/site_extractor/handler.py +++ b/lambdas/site_extractor/handler.py @@ -18,9 +18,14 @@ 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()} @@ -171,10 +176,75 @@ def upsert_site( 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 - skipped = 0 + matched_by_address = 0 + pending = 0 for record in event.get("Records", []): if record["eventName"] not in ("INSERT", "MODIFY"): @@ -187,23 +257,33 @@ def handler(event, context): processed += 1 po = deserialize_image(new_image) po_number = po.get("po_number", "unknown") - - site_code = extract_site_code(po) - if not site_code: - skipped += 1 - logger.info("No site code for PO %s", po_number) - 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 - logger.info("Upserted site %s from PO %s", site_code, po_number) + 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 skipped", + "Batch complete: %d processed, %d extracted (%d via address), %d pending review", processed, extracted, - skipped, + matched_by_address, + pending, ) From 2d79c81e7d714e7a9b8c18f92ad689af507d15f5 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 15:03:43 -0400 Subject: [PATCH 5/5] Update README with pending-site-review table and address fallback docs --- README.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index e11c80d..44ba4e3 100644 --- a/README.md +++ b/README.md @@ -24,7 +24,7 @@ The `po-ingest-site-extractor` Lambda is triggered by the DynamoDB Stream on eve 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 are logged and skipped. +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. @@ -39,6 +39,7 @@ Backfill stats (initial run): 14,825 POs scanned → 9,900 with extractable site - 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`. - **SES:** adds the `PoEmailRule` to the existing `INBOUND_MAIL` receipt rule set (shared with `workorder-ingest`).