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 <noreply@anthropic.com>
This commit is contained in:
Adam Moussa 2026-04-07 12:12:30 -04:00
commit fc690dd958
9 changed files with 709 additions and 0 deletions

12
.gitignore vendored Normal file
View file

@ -0,0 +1,12 @@
__pycache__/
*.py[cod]
*.egg-info/
dist/
build/
.venv/
venv/
node_modules/
cdk.out/
.env
*.eml
lambdas/*/package/

9
cdk/app.py Normal file
View file

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

6
cdk/cdk.json Normal file
View file

@ -0,0 +1,6 @@
{
"app": "python3 app.py",
"context": {
"@aws-cdk/core:bootstrapQualifier": "hnb659fds"
}
}

2
cdk/requirements.txt Normal file
View file

@ -0,0 +1,2 @@
aws-cdk-lib>=2.150.0
constructs>=10.0.0

113
cdk/stack.py Normal file
View file

@ -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")

View file

@ -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"}

View file

@ -0,0 +1,2 @@
anthropic>=0.42.0
boto3>=1.35.0

252
lambdas/web_ui/handler.py Normal file
View file

@ -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'<span style="background:{color};color:#fff;padding:2px 10px;border-radius:12px;font-size:12px;font-weight:500;">{label}</span>'
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", "<br>".join(ship_parts)))
view_url = po.get("view_order_url")
if view_url:
fields.append(("Coupa Link", f'<a href="{view_url}" target="_blank" style="color:#3b82f6;">View in Coupa</a>'))
details_html = ""
for label, value in fields:
if value:
details_html += f"""
<div style="display:flex;padding:8px 0;border-bottom:1px solid #f1f5f9;">
<div style="width:140px;font-size:13px;color:#64748b;font-weight:500;">{label}</div>
<div style="flex:1;font-size:14px;color:#1e293b;">{value}</div>
</div>"""
# Line items
line_items = po.get("line_items") or []
items_html = ""
if line_items:
rows = ""
for item in line_items:
rows += f"""
<tr style="border-bottom:1px solid #f1f5f9;">
<td style="padding:10px;font-size:14px;">{item.get('description', '')}</td>
<td style="padding:10px;font-size:14px;text-align:right;">{fmt_currency(item.get('amount'))}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{item.get('need_by', '') or ''}</td>
<td style="padding:10px;font-size:13px;color:#64748b;">{item.get('category', '') or ''}</td>
</tr>"""
items_html = f"""
<div style="background:#fff;border-radius:10px;padding:24px;box-shadow:0 1px 3px rgba(0,0,0,0.08);">
<h2 style="font-size:16px;margin-bottom:16px;">Line Items ({len(line_items)})</h2>
<table style="width:100%;border-collapse:collapse;">
<thead>
<tr>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Description</th>
<th style="text-align:right;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Amount</th>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Need By</th>
<th style="text-align:left;padding:10px;font-size:12px;text-transform:uppercase;color:#64748b;border-bottom:2px solid #e2e8f0;">Category</th>
</tr>
</thead>
<tbody>{rows}</tbody>
</table>
</div>"""
return f"""<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>PO {po_number} - Sea Haven</title>
<style>
* {{ margin: 0; padding: 0; box-sizing: border-box; }}
body {{ font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', sans-serif; background: #f1f5f9; color: #1e293b; }}
</style>
</head>
<body>
<div style="max-width:800px;margin:0 auto;padding:20px;">
<div style="margin-bottom:20px;">
<a href="/" style="color:#3b82f6;text-decoration:none;font-size:14px;">&larr; All Purchase Orders</a>
</div>
<div style="background:#fff;border-radius:10px;padding:24px;box-shadow:0 1px 3px rgba(0,0,0,0.08);margin-bottom:20px;">
<h1 style="font-size:20px;margin-bottom:16px;">Purchase Order {po_number}</h1>
{details_html}
</div>
{items_html}
</div>
</body>
</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"""
<tr style="border-bottom:1px solid #f1f5f9;cursor:pointer;" onclick="window.location='/po?id={po_number}'">
<td style="padding:12px;font-weight:500;color:#3b82f6;">{po_number}</td>
<td style="padding:12px;">{supplier}</td>
<td style="padding:12px;">{render_badge(email_type, EMAIL_TYPE_COLORS)}</td>
<td style="padding:12px;">{render_badge(status, STATUS_COLORS)}</td>
<td style="padding:12px;text-align:right;">{total}</td>
<td style="padding:12px;color:#64748b;font-size:13px;">{processed}</td>
</tr>"""
return f"""<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Purchase Orders - Sea Haven</title>
<style>
* {{ margin: 0; padding: 0; box-sizing: border-box; }}
body {{ font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', sans-serif; background: #f1f5f9; color: #1e293b; }}
table {{ width: 100%; border-collapse: collapse; }}
tr:hover {{ background: #f8fafc; }}
th {{ text-align: left; padding: 12px; font-size: 12px; text-transform: uppercase; color: #64748b; border-bottom: 2px solid #e2e8f0; }}
</style>
</head>
<body>
<div style="max-width:1100px;margin:0 auto;padding:20px;">
<div style="display:flex;justify-content:space-between;align-items:center;margin-bottom:20px;">
<h1 style="font-size:22px;">Purchase Orders</h1>
<span style="color:#64748b;font-size:14px;">{len(purchase_orders)} most recent</span>
</div>
<div style="background:#fff;border-radius:10px;box-shadow:0 1px 3px rgba(0,0,0,0.08);overflow:hidden;">
<table>
<thead>
<tr>
<th>PO #</th>
<th>Supplier</th>
<th>Type</th>
<th>Status</th>
<th style="text-align:right;">Amount</th>
<th>Processed</th>
</tr>
</thead>
<tbody>
{rows if rows else '<tr><td colspan="6" style="padding:40px;text-align:center;color:#94a3b8;">No purchase orders yet.</td></tr>'}
</tbody>
</table>
</div>
</div>
</body>
</html>"""
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": "<h1>Purchase order not found</h1>",
}
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,
}

80
scripts/reprocess.py Normal file
View file

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