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