From fc690dd958c9d0a2f558cf2917de97a81981c8ae Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Tue, 7 Apr 2026 12:12:30 -0400 Subject: [PATCH] Initial commit: PO email ingestion pipeline CDK stack with SES receipt rule, S3 bucket, email processor Lambda (Claude-powered extraction), web UI Lambda with Function URL, and DynamoDB for storage. Includes reprocessing script for missed emails. Co-Authored-By: Claude Opus 4.6 --- .gitignore | 12 ++ cdk/app.py | 9 + cdk/cdk.json | 6 + cdk/requirements.txt | 2 + cdk/stack.py | 113 ++++++++++ lambdas/email_processor/handler.py | 233 +++++++++++++++++++++ lambdas/email_processor/requirements.txt | 2 + lambdas/web_ui/handler.py | 252 +++++++++++++++++++++++ scripts/reprocess.py | 80 +++++++ 9 files changed, 709 insertions(+) create mode 100644 .gitignore create mode 100644 cdk/app.py create mode 100644 cdk/cdk.json create mode 100644 cdk/requirements.txt create mode 100644 cdk/stack.py create mode 100644 lambdas/email_processor/handler.py create mode 100644 lambdas/email_processor/requirements.txt create mode 100644 lambdas/web_ui/handler.py create mode 100644 scripts/reprocess.py diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5d81fa3 --- /dev/null +++ b/.gitignore @@ -0,0 +1,12 @@ +__pycache__/ +*.py[cod] +*.egg-info/ +dist/ +build/ +.venv/ +venv/ +node_modules/ +cdk.out/ +.env +*.eml +lambdas/*/package/ diff --git a/cdk/app.py b/cdk/app.py new file mode 100644 index 0000000..bd92c9e --- /dev/null +++ b/cdk/app.py @@ -0,0 +1,9 @@ +#!/usr/bin/env python3 +import aws_cdk as cdk +from stack import PoIngestStack + +app = cdk.App() +PoIngestStack(app, "PoIngestStack", + env=cdk.Environment(region="us-east-1"), +) +app.synth() diff --git a/cdk/cdk.json b/cdk/cdk.json new file mode 100644 index 0000000..f28e013 --- /dev/null +++ b/cdk/cdk.json @@ -0,0 +1,6 @@ +{ + "app": "python3 app.py", + "context": { + "@aws-cdk/core:bootstrapQualifier": "hnb659fds" + } +} diff --git a/cdk/requirements.txt b/cdk/requirements.txt new file mode 100644 index 0000000..77978a3 --- /dev/null +++ b/cdk/requirements.txt @@ -0,0 +1,2 @@ +aws-cdk-lib>=2.150.0 +constructs>=10.0.0 diff --git a/cdk/stack.py b/cdk/stack.py new file mode 100644 index 0000000..5e343aa --- /dev/null +++ b/cdk/stack.py @@ -0,0 +1,113 @@ +"""CDK stack for the Coupa PO email ingestion pipeline.""" + +import aws_cdk as cdk +from aws_cdk import ( + Duration, + RemovalPolicy, + Stack, + aws_dynamodb as dynamodb, + aws_lambda as lambda_, + aws_s3 as s3, + aws_s3_notifications as s3n, + aws_ses as ses, + aws_ses_actions as ses_actions, + aws_secretsmanager as secretsmanager, + aws_iam as iam, +) +from constructs import Construct + + +class PoIngestStack(Stack): + def __init__(self, scope: Construct, construct_id: str, **kwargs): + super().__init__(scope, construct_id, **kwargs) + + # --- S3 bucket for raw emails --- + email_bucket = s3.Bucket( + self, "EmailBucket", + bucket_name=f"po-ingest-emails-{self.account}", + removal_policy=RemovalPolicy.RETAIN, + lifecycle_rules=[ + s3.LifecycleRule(expiration=Duration.days(90)), + ], + ) + + # --- 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", + ) + + # --- Secrets Manager for Anthropic API key --- + anthropic_secret = secretsmanager.Secret( + self, "AnthropicApiKey", + secret_name="po-ingest/anthropic-api-key", + description="Anthropic API key for Coupa PO email parsing", + ) + + # --- Lambda function --- + email_processor = lambda_.Function( + self, "EmailProcessor", + function_name="po-email-processor", + runtime=lambda_.Runtime.PYTHON_3_12, + handler="handler.handler", + code=lambda_.Code.from_asset("../lambdas/email_processor/package"), + timeout=Duration.seconds(60), + memory_size=256, + environment={ + "PO_TABLE": "purchase-orders", + "ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn, + }, + ) + + # Grant permissions + email_bucket.grant_read(email_processor) + po_table.grant_read_write_data(email_processor) + anthropic_secret.grant_read(email_processor) + + # S3 event notification → Lambda + email_bucket.add_event_notification( + s3.EventType.OBJECT_CREATED, + s3n.LambdaDestination(email_processor), + s3.NotificationKeyFilter(prefix="inbound/"), + ) + + # --- SES Receipt Rule --- + # Reuse the existing INBOUND_MAIL rule set (shared with workorder-ingest) + rule_set = ses.ReceiptRuleSet.from_receipt_rule_set_name( + self, "ExistingRuleSet", "INBOUND_MAIL", + ) + + rule_set.add_rule( + "PoEmailRule", + recipients=["amazon_po@int.seahaven.com"], + actions=[ + ses_actions.S3( + bucket=email_bucket, + object_key_prefix="inbound/", + ), + ], + ) + + # --- Web UI Lambda --- + web_ui = lambda_.Function( + self, "WebUI", + function_name="po-web-ui", + runtime=lambda_.Runtime.PYTHON_3_12, + handler="handler.handler", + code=lambda_.Code.from_asset("../lambdas/web_ui"), + timeout=Duration.seconds(60), + memory_size=256, + environment={ + "PO_TABLE": "purchase-orders", + }, + ) + + po_table.grant_read_data(web_ui) + + # Function URL for direct access + web_url = web_ui.add_function_url( + auth_type=lambda_.FunctionUrlAuthType.NONE, + ) + + cdk.CfnOutput(self, "WebUIUrl", value=web_url.url, description="PO Dashboard URL") diff --git a/lambdas/email_processor/handler.py b/lambdas/email_processor/handler.py new file mode 100644 index 0000000..4f28086 --- /dev/null +++ b/lambdas/email_processor/handler.py @@ -0,0 +1,233 @@ +""" +PO email processor Lambda. + +Triggered by S3 events when SES delivers a Coupa PO email. +Parses the raw email, sends it to Claude for structured extraction, +then writes the result to the purchase-orders DynamoDB table. +""" + +import email +import json +import logging +import os +import re +from datetime import datetime +from decimal import Decimal +from email import policy + +import anthropic +import boto3 + +logger = logging.getLogger() +logger.setLevel(logging.INFO) + +s3 = boto3.client("s3") +dynamodb = boto3.resource("dynamodb") + +PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders") +ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN") + +EXTRACTION_PROMPT = """\ +You are an email parser for a purchase order system. +The emails come from Coupa (a procurement platform) and contain purchase order +notifications from Amazon. + +Analyze the following email and extract structured data. Return ONLY valid JSON with these fields: + +{ + "email_type": "new_po" | "cancellation", + "po_number": "string or null", + "po_status": "string or null", + "source_system": "coupa", + "submitted_by": "string or null", + "on_behalf_of": "string or null", + "order_date": "string or null", + "revision_date": "string or null", + "last_opened": "string or null", + "acknowledged_at": "string or null", + "payment_terms": "string or null", + "requisition_number": "string or null", + "department": "string or null", + "view_order_url": "URL string or null", + "supplier": { + "name": "string or null" + }, + "ship_to": { + "name": "string or null", + "address": "string or null", + "location_code": "string or null", + "attn": "string or null" + }, + "total_amount": 0.0, + "currency": "USD", + "line_items": [ + { + "description": "string", + "amount": 0.0, + "currency": "USD", + "need_by": "date string or null", + "category": "string or null", + "account_code": "string or null", + "period": "string or null" + } + ] +} + +Rules: +- "email_type" detection: + - "new_po": email announces a new or revised purchase order + - "cancellation": email announces a PO has been cancelled +- 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 +- 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 +""" + + +def get_anthropic_client() -> anthropic.Anthropic: + """Create Anthropic client, fetching API key from Secrets Manager.""" + if ANTHROPIC_API_KEY_SECRET_ARN: + secrets = boto3.client("secretsmanager") + secret = secrets.get_secret_value(SecretId=ANTHROPIC_API_KEY_SECRET_ARN) + api_key = secret["SecretString"] + return anthropic.Anthropic(api_key=api_key) + # Fall back to ANTHROPIC_API_KEY env var (for local testing) + return anthropic.Anthropic() + + +def parse_raw_email(raw_bytes: bytes) -> dict: + """Parse a raw MIME email into subject, sender, and body text.""" + msg = email.message_from_bytes(raw_bytes, policy=policy.default) + + subject = msg.get("Subject", "") + sender = msg.get("From", "") + to = msg.get("To", "") + date = msg.get("Date", "") + + body = "" + if msg.is_multipart(): + for part in msg.walk(): + content_type = part.get_content_type() + if content_type == "text/plain": + body = part.get_content() + break + elif content_type == "text/html" and not body: + body = part.get_content() + else: + body = msg.get_content() + + return { + "subject": subject, + "sender": sender, + "to": to, + "date": date, + "body": body, + } + + +def extract_with_claude(email_data: dict) -> dict: + """Send parsed email to Claude for structured extraction.""" + client = get_anthropic_client() + + email_text = ( + f"Subject: {email_data['subject']}\n" + f"From: {email_data['sender']}\n" + f"To: {email_data['to']}\n" + f"Date: {email_data['date']}\n" + f"\n---\n\n" + f"{email_data['body']}" + ) + + response = client.messages.create( + model="claude-haiku-4-5-20251001", + max_tokens=2048, + messages=[ + { + "role": "user", + "content": f"{EXTRACTION_PROMPT}\n\nEMAIL:\n{email_text}", + } + ], + ) + + response_text = response.content[0].text + + # Extract JSON from response (handle markdown code blocks) + json_match = re.search(r"```(?:json)?\s*(.*?)```", response_text, re.DOTALL) + if json_match: + response_text = json_match.group(1) + + return json.loads(response_text.strip(), parse_float=Decimal) + + +def save_new_po(parsed: dict, s3_key: str): + """Insert a new PO into DynamoDB. Skips if po_number already exists.""" + table = dynamodb.Table(PO_TABLE) + now = datetime.utcnow().isoformat() + + # Add metadata fields + parsed["raw_s3_key"] = s3_key + parsed["processed_at"] = now + + # Build item, stripping None values + item = {k: v for k, v in parsed.items() if v is not None} + + try: + table.put_item( + Item=item, + ConditionExpression="attribute_not_exists(po_number)", + ) + logger.info(f"Created PO {parsed['po_number']}") + except dynamodb.meta.client.exceptions.ConditionalCheckFailedException: + logger.info(f"PO {parsed['po_number']} already exists, skipping insert") + + +def save_cancellation(parsed: dict, s3_key: str): + """Update an existing PO's status to Cancelled.""" + table = dynamodb.Table(PO_TABLE) + + table.update_item( + Key={"po_number": parsed["po_number"]}, + UpdateExpression="SET po_status = :status, cancelled_at = :cancelled_at, raw_s3_key = :s3_key", + ExpressionAttributeValues={ + ":status": "Cancelled", + ":cancelled_at": datetime.utcnow().isoformat(), + ":s3_key": s3_key, + }, + ) + logger.info(f"Cancelled PO {parsed['po_number']}") + + +def handler(event, context): + """Lambda entry point. Triggered by S3 ObjectCreated events.""" + for record in event.get("Records", []): + bucket = record["s3"]["bucket"]["name"] + key = record["s3"]["object"]["key"] + s3_key = f"s3://{bucket}/{key}" + + logger.info(f"Processing email: {s3_key}") + + # Fetch raw email from S3 + response = s3.get_object(Bucket=bucket, Key=key) + raw_email = response["Body"].read() + + # Parse the raw email + email_data = parse_raw_email(raw_email) + logger.info(f"Subject: {email_data['subject']}") + + # Extract structured data with Claude + parsed = extract_with_claude(email_data) + logger.info(f"Parsed: type={parsed.get('email_type')}, po={parsed.get('po_number')}") + + if not parsed.get("po_number"): + logger.warning(f"No PO number found in email, skipping: {key}") + continue + + # Route by email type + if parsed.get("email_type") == "cancellation": + save_cancellation(parsed, s3_key) + else: + save_new_po(parsed, s3_key) + + return {"statusCode": 200, "body": "OK"} diff --git a/lambdas/email_processor/requirements.txt b/lambdas/email_processor/requirements.txt new file mode 100644 index 0000000..acef34e --- /dev/null +++ b/lambdas/email_processor/requirements.txt @@ -0,0 +1,2 @@ +anthropic>=0.42.0 +boto3>=1.35.0 diff --git a/lambdas/web_ui/handler.py b/lambdas/web_ui/handler.py new file mode 100644 index 0000000..22d6d44 --- /dev/null +++ b/lambdas/web_ui/handler.py @@ -0,0 +1,252 @@ +""" +Web UI Lambda. + +Serves a simple HTML dashboard for viewing purchase orders. +Accessed via Lambda Function URL. +""" + +import json +import os +from decimal import Decimal + +import boto3 + +dynamodb = boto3.resource("dynamodb") + +PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders") + + +def get_purchase_orders(limit=500): + table = dynamodb.Table(PO_TABLE) + items = [] + response = table.scan() + items.extend(response.get("Items", [])) + while "LastEvaluatedKey" in response: + response = table.scan(ExclusiveStartKey=response["LastEvaluatedKey"]) + items.extend(response.get("Items", [])) + items.sort(key=lambda x: x.get("processed_at", ""), reverse=True) + return items[:limit] + + +def render_badge(value, color_map): + if not value: + value = "unknown" + color = color_map.get(value.lower(), "#9ca3af") + label = value.replace("_", " ").title() + return f'{label}' + + +STATUS_COLORS = { + "issued": "#3b82f6", + "pending buyer action": "#f59e0b", + "open": "#3b82f6", + "closed": "#10b981", + "cancelled": "#ef4444", + "soft closed": "#6b7280", +} + +EMAIL_TYPE_COLORS = { + "new_po": "#3b82f6", + "cancellation": "#ef4444", +} + + +def fmt_currency(val): + if val is None: + return "" + if isinstance(val, Decimal): + val = float(val) + if isinstance(val, (int, float)): + return f"${val:,.2f}" + return str(val) + + +def render_po_detail(po): + po_number = po.get("po_number", "") + + fields = [ + ("PO Number", po_number), + ("Status", render_badge(po.get("po_status", ""), STATUS_COLORS)), + ("Email Type", render_badge(po.get("email_type", ""), EMAIL_TYPE_COLORS)), + ("Total Amount", fmt_currency(po.get("total_amount"))), + ("Currency", po.get("currency")), + ("Supplier", (po.get("supplier") or {}).get("name")), + ("Submitted By", po.get("submitted_by")), + ("On Behalf Of", po.get("on_behalf_of")), + ("Order Date", po.get("order_date")), + ("Revision Date", po.get("revision_date")), + ("Payment Terms", po.get("payment_terms")), + ("Requisition #", po.get("requisition_number")), + ("Department", po.get("department")), + ("Processed At", po.get("processed_at")), + ] + + ship_to = po.get("ship_to") or {} + if any(ship_to.values()): + ship_parts = [] + if ship_to.get("name"): + ship_parts.append(ship_to["name"]) + if ship_to.get("address"): + ship_parts.append(ship_to["address"]) + if ship_to.get("location_code"): + ship_parts.append(f"Location: {ship_to['location_code']}") + if ship_to.get("attn"): + ship_parts.append(f"Attn: {ship_to['attn']}") + fields.append(("Ship To", "
".join(ship_parts))) + + view_url = po.get("view_order_url") + if view_url: + fields.append(("Coupa Link", f'View in Coupa')) + + details_html = "" + for label, value in fields: + if value: + details_html += f""" +
+
{label}
+
{value}
+
""" + + # Line items + line_items = po.get("line_items") or [] + items_html = "" + if line_items: + rows = "" + for item in line_items: + rows += f""" + + {item.get('description', '')} + {fmt_currency(item.get('amount'))} + {item.get('need_by', '') or ''} + {item.get('category', '') or ''} + """ + + items_html = f""" +
+

Line Items ({len(line_items)})

+ + + + + + + + + + {rows} +
DescriptionAmountNeed ByCategory
+
""" + + return f""" + + + + + PO {po_number} - Sea Haven + + + +
+
+ ← All Purchase Orders +
+
+

Purchase Order {po_number}

+ {details_html} +
+ {items_html} +
+ +""" + + +def render_po_list(purchase_orders): + rows = "" + for po in purchase_orders: + po_number = po.get("po_number", "") + supplier = (po.get("supplier") or {}).get("name", "") + status = po.get("po_status", "") + email_type = po.get("email_type", "") + total = fmt_currency(po.get("total_amount")) + processed = (po.get("processed_at") or "")[:16] + + rows += f""" + + {po_number} + {supplier} + {render_badge(email_type, EMAIL_TYPE_COLORS)} + {render_badge(status, STATUS_COLORS)} + {total} + {processed} + """ + + return f""" + + + + + Purchase Orders - Sea Haven + + + +
+
+

Purchase Orders

+ {len(purchase_orders)} most recent +
+
+ + + + + + + + + + + + + {rows if rows else ''} + +
PO #SupplierTypeStatusAmountProcessed
No purchase orders yet.
+
+
+ +""" + + +def handler(event, context): + path = event.get("rawPath", "/") + qs = event.get("queryStringParameters") or {} + + if path == "/po" and "id" in qs: + po_number = qs["id"] + table = dynamodb.Table(PO_TABLE) + result = table.get_item(Key={"po_number": po_number}) + po = result.get("Item") + if not po: + return { + "statusCode": 404, + "headers": {"Content-Type": "text/html"}, + "body": "

Purchase order not found

", + } + html = render_po_detail(po) + else: + purchase_orders = get_purchase_orders() + html = render_po_list(purchase_orders) + + return { + "statusCode": 200, + "headers": {"Content-Type": "text/html"}, + "body": html, + } diff --git a/scripts/reprocess.py b/scripts/reprocess.py new file mode 100644 index 0000000..6535be5 --- /dev/null +++ b/scripts/reprocess.py @@ -0,0 +1,80 @@ +""" +Re-invoke the po-email-processor Lambda for every object still +sitting under the inbound/ prefix in the email bucket. + +Usage: + python scripts/reprocess.py # dry-run (list only) + python scripts/reprocess.py --execute # actually invoke +""" + +import argparse +import json + +import boto3 + +sts = boto3.client("sts") +s3 = boto3.client("s3") +lambda_client = boto3.client("lambda") + +FUNCTION_NAME = "po-email-processor" + + +def get_bucket_name() -> str: + account_id = sts.get_caller_identity()["Account"] + return f"po-ingest-emails-{account_id}" + + +def list_inbound_keys(bucket: str) -> list[str]: + keys = [] + paginator = s3.get_paginator("list_objects_v2") + for page in paginator.paginate(Bucket=bucket, Prefix="inbound/"): + for obj in page.get("Contents", []): + keys.append(obj["Key"]) + return keys + + +def build_s3_event(bucket: str, key: str) -> dict: + return { + "Records": [ + { + "s3": { + "bucket": {"name": bucket}, + "object": {"key": key}, + } + } + ] + } + + +def main(): + parser = argparse.ArgumentParser(description="Reprocess missed PO emails") + parser.add_argument("--execute", action="store_true", help="Actually invoke the Lambda (default is dry-run)") + args = parser.parse_args() + + bucket = get_bucket_name() + keys = list_inbound_keys(bucket) + + if not keys: + print("No objects found under inbound/ — nothing to reprocess.") + return + + print(f"Found {len(keys)} email(s) in s3://{bucket}/inbound/\n") + + for key in keys: + if args.execute: + print(f" Invoking for {key} ... ", end="", flush=True) + resp = lambda_client.invoke( + FunctionName=FUNCTION_NAME, + InvocationType="Event", # async — don't wait for each one + Payload=json.dumps(build_s3_event(bucket, key)), + ) + print(f"status {resp['StatusCode']}") + else: + print(f" [dry-run] {key}") + + if not args.execute: + print(f"\nDry run complete. Re-run with --execute to invoke the Lambda.") + + +if __name__ == "__main__": + main()