mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-10-01 01:53:12 +00:00
The From header and any raw-MIME Authentication-Results copies are attacker-forgeable, so a forged email to apm@int.seahaven.com or amazon_po@int.seahaven.com could create or mutate a WO/PO (INFRA-107, CRITICAL). Both S3-triggered email processors now authenticate the sender against the Authentication-Results header SES itself prepends at delivery: only the topmost header is consulted, its authserv-id must be amazonses.com, and it must carry dkim=pass for a domain in the per-pipeline ALLOWED_DKIM_DOMAINS env var (comma-separated, set in CDK so ops can adjust without code changes). Allowlists come from live traffic observed 2026-07-15 on both ingest buckets: WO mail arrives via the apm@ Google Groups forward, which re-signs as seahaven.com (the hxgnsmartcloud.com signature does not survive the forward); PO mail passes for amazon.coupahost.com. amazonses.com also passes on PO mail but is deliberately excluded -- every SES customer's outbound mail passes for it. Every failure path (env var unset, header missing or unparseable, verdict fail, unaligned domain) rejects the email: a structured warning with the reason and S3 key is logged and the record skipped without erroring the invocation, so rejected mail causes no Lambda retries or DLQ messages. Handler signatures and event sources are unchanged. Refs: INFRA-107
278 lines
9.4 KiB
Python
278 lines
9.4 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 json
|
|
import logging
|
|
import os
|
|
import re
|
|
from datetime import datetime
|
|
from email import policy
|
|
|
|
import anthropic
|
|
import boto3
|
|
from ses_auth import authenticate_inbound_email
|
|
|
|
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")
|
|
|
|
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 if configured."""
|
|
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 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,
|
|
}
|
|
|
|
|
|
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"]
|
|
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()
|
|
|
|
# Fail-closed sender authentication (INFRA-107): only mail with an
|
|
# SES-stamped dkim=pass verdict for an allowlisted domain may create
|
|
# or update work orders. Rejected mail is logged and skipped without
|
|
# erroring the invocation (no retries / DLQ spam).
|
|
if not authenticate_inbound_email(raw_email, s3_key):
|
|
continue
|
|
|
|
# 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')}, 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
|
|
|
|
# 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"}
|