procurement-ingest/lambdas/wo/email_processor/handler.py
Adam Moussa f175323939 Validate sender and require Secrets Manager key
The PO and WO email processors acted on email from any sender — the
public addresses (amazon_po@, apm@) accept mail from anyone, so an
attacker could forge POs/WOs. Reject email whose verified sender
domain is not on a configurable allowlist (ALLOWED_SENDER_DOMAINS),
and drop messages with an explicit SES SPF/DKIM/spam/virus failure.

Also remove the silent fallback to a plaintext ANTHROPIC_API_KEY env
var; require ANTHROPIC_API_KEY_SECRET_ARN and raise if absent so a
misconfigured deploy fails loudly instead of using an unmanaged key.
2026-06-17 11:37:29 -04:00

356 lines
12 KiB
Python

"""
Email processor Lambda.
Triggered by S3 events when SES delivers an email.
Parses the raw email, sends it to Claude for structured extraction,
then writes the result to DynamoDB.
"""
import email
import email.utils
import json
import logging
import os
import re
from datetime import datetime
from email import policy
import anthropic
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
s3 = boto3.client("s3")
dynamodb = boto3.resource("dynamodb")
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN")
# Allowlist of sender domains permitted to create/modify work orders. APM emails
# originate from Amazon's APM (via Hexagon EAM / HxGN SmartCloud). The leading dot
# semantics ("domain or subdomain") are applied in is_sender_allowed. Override via
# the ALLOWED_SENDER_DOMAINS env var (comma-separated) without a code redeploy. An
# attacker who can email apm@int.seahaven.com but cannot forge a verified sender in
# this list is rejected before any DynamoDB write.
DEFAULT_ALLOWED_SENDER_DOMAINS = "amazon.com,hxgnsmartcloud.com,hexagon.com"
ALLOWED_SENDER_DOMAINS = [
d.strip().lower().lstrip("@")
for d in os.environ.get(
"ALLOWED_SENDER_DOMAINS", DEFAULT_ALLOWED_SENDER_DOMAINS
).split(",")
if d.strip()
]
EXTRACTION_PROMPT = """\
You are an email parser for a facilities maintenance work order system.
The emails come from Amazon's APM system (via Hexagon EAM / HxGN SmartCloud).
Analyze the following email and extract structured data. Return ONLY valid JSON with these fields:
{
"email_type": "new_work_order" | "update" | "comment" | "cancellation",
"work_order_id": "string or null",
"description": "work order description or null",
"status": "new" | "assigned" | "in_progress" | "on_hold" | "completed" | "cancelled" | "unknown",
"site_code": "building/site code like WIL1, ZDL8, etc. or null",
"building": "full building identifier or null",
"address": "physical address or null",
"severity": "severity level or null",
"priority": "priority level or null",
"date_reported": "ISO 8601 date or null",
"scheduled_start": "ISO 8601 date or null",
"due_date": "ISO 8601 date or null",
"assigned_to": "person/team assigned or null",
"commenter": "person who left a comment or null",
"comment_text": "the comment text or null",
"comment_time": "ISO 8601 datetime of the comment or null"
}
Rules:
- "email_type" detection:
- "new_work_order": email announces a new WO assignment
- "comment": email contains a new comment on an existing WO
- "cancellation": email announces a WO has been cancelled
- "update": any other update to an existing WO (status change, reassignment, etc.)
- Extract the site_code from the building field (e.g., "WIL1" from "building WIL1")
- Dates should be converted to ISO 8601 format
- 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.
Requires ANTHROPIC_API_KEY_SECRET_ARN to be set. The previous silent
fallback to a plaintext ANTHROPIC_API_KEY env var is removed: a misconfigured
deploy must fail loudly rather than quietly run on an unmanaged key.
"""
if not ANTHROPIC_API_KEY_SECRET_ARN:
raise RuntimeError(
"ANTHROPIC_API_KEY_SECRET_ARN is not set; refusing to fall back to a "
"plaintext API key. Configure the Secrets Manager 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)
def extract_sender_domain(sender: str) -> str | None:
"""Extract the lowercased domain from a From header value."""
if not sender:
return None
_, addr = email.utils.parseaddr(sender)
if "@" not in addr:
return None
return addr.rsplit("@", 1)[1].strip().lower()
def is_sender_allowed(sender: str) -> bool:
"""Return True if the sender domain is in the configured allowlist.
A domain matches if it equals an allowlist entry or is a subdomain of one.
"""
domain = extract_sender_domain(sender)
if not domain:
return False
for allowed in ALLOWED_SENDER_DOMAINS:
if domain == allowed or domain.endswith("." + allowed):
return True
return False
def ses_auth_failed(email_data: dict) -> bool:
"""Return True if SES recorded a hard SPF/DKIM/spam/virus failure.
Only blocks on an explicit "fail" so mails delivered before receipt-rule
scanning was enabled (no header) are not silently dropped; the domain
allowlist remains the primary gate.
"""
spam = (email_data.get("ses_spam_verdict") or "").upper()
virus = (email_data.get("ses_virus_verdict") or "").upper()
if spam == "FAIL" or virus == "FAIL":
return True
auth = (email_data.get("authentication_results") or "").lower()
if "spf=fail" in auth or "dkim=fail" in auth:
return True
return False
def parse_raw_email(raw_bytes: bytes) -> dict:
"""Parse a raw email into subject, sender, body text."""
msg = email.message_from_bytes(raw_bytes, policy=policy.default)
subject = msg.get("Subject", "")
sender = msg.get("From", "")
to = msg.get("To", "")
cc = msg.get("Cc", "")
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,
"cc": cc,
"date": date,
"body": body,
# SES stamps these on the stored object when receipt-rule scanning is on.
"ses_spam_verdict": msg.get("X-SES-Spam-Verdict", ""),
"ses_virus_verdict": msg.get("X-SES-Virus-Verdict", ""),
"authentication_results": msg.get("Authentication-Results", ""),
}
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"CC: {email_data['cc']}\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=1024,
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())
def save_work_order(parsed: dict, s3_key: str):
"""Create or update a work order in DynamoDB."""
table = dynamodb.Table(WORK_ORDERS_TABLE)
work_order_id = parsed["work_order_id"]
now = datetime.utcnow().isoformat()
# Build update expression dynamically from non-null fields
field_map = {
"description": "description",
"status": "wo_status", # 'status' is a DynamoDB reserved word
"site_code": "site_code",
"building": "building",
"address": "address",
"severity": "severity",
"priority": "priority",
"date_reported": "date_reported",
"scheduled_start": "scheduled_start",
"due_date": "due_date",
"assigned_to": "assigned_to",
}
update_parts = ["#updated_at = :updated_at", "#source_key = :source_key"]
attr_names = {
"#updated_at": "updated_at",
"#source_key": "source_email_s3_key",
}
attr_values = {
":updated_at": now,
":source_key": s3_key,
}
for src_field, dynamo_field in field_map.items():
value = parsed.get(src_field)
if value is not None:
placeholder = f":{dynamo_field}"
name_placeholder = f"#{dynamo_field}"
update_parts.append(f"{name_placeholder} = {placeholder}")
attr_names[name_placeholder] = dynamo_field
attr_values[placeholder] = value
# For new items, set created_at
update_parts.append("#created_at = if_not_exists(#created_at, :created_at)")
attr_names["#created_at"] = "created_at"
attr_values[":created_at"] = now
# Customer is always AMAZON for now
update_parts.append("#customer = :customer")
attr_names["#customer"] = "customer"
attr_values[":customer"] = "AMAZON"
# Track the record type (new_work_order, update, comment)
email_type = parsed.get("email_type")
if email_type:
update_parts.append("#record_type = :record_type")
attr_names["#record_type"] = "record_type"
attr_values[":record_type"] = email_type
table.update_item(
Key={"work_order_id": work_order_id},
UpdateExpression="SET " + ", ".join(update_parts),
ExpressionAttributeNames=attr_names,
ExpressionAttributeValues=attr_values,
)
logger.info(f"Saved work order {work_order_id}")
def save_event(parsed: dict, s3_key: str):
"""Save an event to the events table. Every email creates an event entry."""
table = dynamodb.Table(COMMENTS_TABLE)
work_order_id = parsed["work_order_id"]
email_type = parsed.get("email_type", "unknown")
event_time = parsed.get("comment_time") or datetime.utcnow().isoformat()
event_id = f"{work_order_id}#{event_time}"
item = {
"work_order_id": work_order_id,
"comment_id": event_id, # keeping key name for table compatibility
"record_type": email_type,
"commenter": parsed.get("commenter") or "",
"text": parsed.get("comment_text") or "",
"created_at": event_time,
"source_email_s3_key": s3_key,
"ingested_at": datetime.utcnow().isoformat(),
}
table.put_item(Item=item)
logger.info(f"Saved event {event_id} (type={email_type})")
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"]
logger.info(f"Processing email: s3://{bucket}/{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']}")
# Authz: only act on email from an allowlisted sender domain that passed
# SES auth checks. Anyone can email apm@int.seahaven.com, but only
# legitimate APM/Hexagon senders may create or mutate work orders.
sender = email_data.get("sender", "")
if not is_sender_allowed(sender):
logger.warning(
f"Rejecting email from disallowed sender '{sender}' "
f"(domain not in allowlist): {key}"
)
continue
if ses_auth_failed(email_data):
logger.warning(
f"Rejecting email from '{sender}' due to SES SPF/DKIM/spam "
f"failure: {key}"
)
continue
# Extract structured data with Claude
parsed = extract_with_claude(email_data)
logger.info(
f"Parsed: type={parsed.get('email_type')}, wo={parsed.get('work_order_id')}"
)
if not parsed.get("work_order_id"):
logger.warning(f"No work order ID found in email, skipping: {key}")
continue
s3_key = f"s3://{bucket}/{key}"
# Always upsert the work order with any new info
save_work_order(parsed, s3_key)
# Save every email as an event for history tracking
save_event(parsed, s3_key)
return {"statusCode": 200, "body": "OK"}