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