procurement-ingest/lambdas/email_processor/handler.py

246 lines
7.8 KiB
Python
Raw Normal View History

"""
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"
},
"site_code": "string or null",
"ship_to": {
"name": "string or null",
"address": "string or null",
"street": "string or null",
"city": "string or null",
"state": "string or null",
"zip": "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
- Ship-to street, city, state, and zip should be parsed from the address into separate fields
- "site_code" is the Amazon facility code (e.g., "SNY5", "DFW6", "WND1") — a 3-5 character alphanumeric code identifying the delivery site. Look for it in:
- The ship-to name, e.g., "Amazon.com Services LLC - SNY5" or "Amazon.com Services LLC (WFB1)"
- The ATTN line, e.g., "ATTN: Wagon Wheel DS - WTN1"
- Line item descriptions, e.g., "WND1 - 2024 - Plumbing PM"
- Anywhere else in the email where a facility code appears
- If the ship-to name IS the site code (e.g., just "DBU2"), use that
- 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"}