mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 15:23:13 +00:00
272 lines
9 KiB
Python
272 lines
9 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
|
||
|
|
|
||
|
|
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"]
|
||
|
|
|
||
|
|
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']}")
|
||
|
|
|
||
|
|
# 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"}
|