mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 04:53:12 +00:00
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.
This commit is contained in:
parent
031a0e1c60
commit
ec416079f5
3 changed files with 342 additions and 5 deletions
78
cdk/stack.py
78
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",
|
||||
)
|
||||
|
|
|
|||
193
lambdas/site_extractor/handler.py
Normal file
193
lambdas/site_extractor/handler.py
Normal file
|
|
@ -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,
|
||||
)
|
||||
76
scripts/backfill_sites.py
Normal file
76
scripts/backfill_sites.py
Normal file
|
|
@ -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()
|
||||
Loading…
Add table
Reference in a new issue