procurement-ingest/lambdas/wo/email_processor/handler.py

279 lines
9.3 KiB
Python
Raw Normal View History

Merge workorder-ingest into unified procurement repo (#22) * Merge workorder-ingest pipeline into unified repo Move PO lambdas under lambdas/po/, add WO pipeline under lambdas/wo/. Two independent CloudFormation stacks in one CDK app. Fix WO stack compliance: ARM64 architecture, 60-day log retention, aarch64 bundling, RETAIN on Anthropic secret. Remove stale CodePipeline buildspec. * Fix test_local.py import path and remove dead shared/models.py test_local.py referenced the old lambdas/email_processor path. Updated to lambdas/wo/email_processor. Removed shared/ directory entirely as nothing imports from it. * Escape HTML in both web UI dashboards to prevent XSS Both Function URLs are public (auth_type=NONE) and render email-derived content via f-strings. Attacker-crafted emails could inject scripts. Added html.escape() on all interpolated values in both PO and WO dashboards. * Add pagination to WO web UI scan get_work_orders() only fetched the first 1MB page from DynamoDB. Loop on LastEvaluatedKey to match the PO web UI pattern. * Fix esc(None) TypeError and javascript: scheme in PO web UI Coerce supplier name through `or ""` before escaping to handle nested None from DynamoDB. Add scheme allowlist on view_order_url to block javascript:/data: hrefs from LLM-extracted URLs. * Fix WO render_badge None guard, updated_at slice, and backfill path Add null guard to WO render_badge matching the PO version. Use `or ""` before slicing updated_at to handle explicit None values. Fix backfill_sites.py sys.path to use new lambdas/po/site_extractor. * Harden WO web UI and fix JS-context XSS in both dashboards - Use json.dumps for onclick URLs to prevent JS string breakout - Add .lower() to WO render_badge color lookup matching PO pattern - Add pagination to get_comments query - Cap get_work_orders to 500 results matching PO pattern * Apply ruff formatting to web UI handlers
2026-05-12 15:21:06 -04:00
"""
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:
Land safe fixes from 2026-06-17 security sweep (#97) * Remove gratuitous KMS grant on shared DynamoDB CMK wo-email-processor held grant_encrypt_decrypt on the shared seahaven-dynamodb CMK, but the WorkOrders/WorkOrderComments tables are not encrypted with that CMK. The grant was dead weight that extended the WO processor's decrypt reach to the CMK protecting the purchase-orders table (cross-stack decrypt). Drop it to restore least privilege; re-add as part of the table CMK migration (INFRA-6). Refs: INFRA-6 * Require Secrets Manager key for Anthropic client Remove the silent fallback to a plaintext ANTHROPIC_API_KEY env var in both email processors; require ANTHROPIC_API_KEY_SECRET_ARN and raise if absent so a misconfigured deploy fails loudly instead of using an unmanaged key. Adapted from f175323 on security/sweep-2026-06-17. The From-header sender-domain allowlist from that commit is intentionally dropped: the From header is spoofable (INFRA-107, confirmed critical) and sender authentication is being reworked in a separate PR. Refs: INFRA-107 * Merge PO revisions and handle out-of-order events save_revision did a full put_item overwrite, so a revision omitting line_items/supplier permanently deleted them. save_new_po used a conditional put that silently dropped the PO when an out-of-order cancellation had already created a skeleton row. Switch both to field-level merge update_items: a revision now SETs only the fields it carries, and a new_po backfills data into a pre-existing Cancelled skeleton while preserving the Cancelled status. No email can now delete data established by an earlier one. * Gate web UIs behind auth and escape currency XSS The po-web-ui and workorder-web-ui handlers had no auth: any invocation path returned the full PO/WO DB. Add a fail-closed shared-secret gate (X-Auth-Token / Bearer, constant-time compared to WEB_UI_AUTH_TOKEN) so a future re-attached Function URL cannot re-expose the data (URLs removed under INFRA-74). Wire the token from the SSM String param /procurement-ingest/web-ui-auth-token. Also fix stored XSS in po-web-ui fmt_currency: the non-numeric fallback returned str(val) unescaped, so a prompt-injected email could make Claude emit total_amount as <script>. Escape it. Refs: INFRA-74 * Document sweep security fixes and merge semantics Update the README for the 2026-06-17 security sweep: required Secrets Manager key (no plaintext env fallback), web UI auth gate + SSM token setup step, output-escaping note, and the new PO revision/cancellation merge behavior. Adapted from d91f45e on security/sweep-2026-06-17; the sender allowlist documentation is dropped along with the allowlist itself (deferred to the INFRA-107 sender-authentication rework). Refs: INFRA-107 * fix: resolve web UI auth token from Secrets Manager at runtime Replace the plaintext SSM String parameter with a Secrets Manager secret referenced by ARN only. The token is fetched and cached at module level on first invocation, keeping shared secrets out of CloudFormation templates and Lambda environment variables. Refs: PR-97 * Add TTL to web UI auth token cache for rotation The web-ui handlers cached the Secrets Manager auth token at module level with no expiry, so a rotated secret was only picked up when the warm container recycled — an emergency rotation could take hours to take effect. Cache the fetched value for a 5-minute TTL instead, so a rotated token propagates within the TTL while still avoiding a Secrets Manager call on every request. Still fails closed when the secret is unset or unreadable. Refs: INFRA-74 * Log Secrets Manager failures in web UI auth token fetch The web UI auth gate correctly fails closed when the shared token cannot be read, but _get_auth_token() swallowed every exception silently. A Secrets Manager permission or config error then made every request 401 with no operational signal, leaving an outage indistinguishable from ordinary unauthenticated traffic. Add a module-level logger to both web_ui handlers and log the fetch failure with logger.exception() in the except block before returning None. Behavior is unchanged (still fails closed); the failure is now visible in CloudWatch. The secret value is never logged. The two handlers stay byte-consistent in the mirrored _get_auth_token() region. The companion finding on the CDK import of the shared procurement-ingest/web-ui-auth-token secret was evaluated and left as-is: the token is a single secret shared by both the PO and WO stacks, so from_secret_name_v2 (which scopes grant_read via the standard 6-char suffix wildcard) is correct; making it a CDK-managed Secret in both stacks would collide the two stacks on the same explicit secret name at deploy time. Refs: INFRA-74 * Make Cancelled PO status sticky via atomic write The PO merge path read status with a get_item (_is_cancelled) and then wrote with an unconditional update_item. Two defects followed from this: - Race (Issue A): a cancellation landing between the read and the write was silently un-cancelled by a revision carrying a non-cancelled po_status — a TOCTOU on a table with concurrent email processing. - Over-broad strip (Issue B): save_revision dropped po_status whenever the PO was Cancelled, so legitimate status updates on non-cancelled POs and status-less revisions were affected rather than only the true un-cancel transition. Enforce the invariant server-side instead. "Cancelled" is a sticky, authoritative status: once set, later new_po/revision emails may enrich other fields but must never move it to a non-cancelled status. When the payload carries a non-cancelled po_status, _merge_update issues the update_item guarded by ConditionExpression "attribute_not_exists(po_status) OR po_status <> :marker", evaluated atomically at write time, so a cancellation that lands first always wins. On ConditionalCheckFailedException the same fields are re-written without po_status/cancelled_at, enriching the record while Cancelled sticks. Payloads with no status change, or an already -Cancelled status, take a plain merge — the status is only ever suppressed on a real un-cancel. This removes the non-atomic get_item from the write path; _is_cancelled is deleted. Key schema and attribute names are unchanged, so the cross-stack purchase-orders contract (read-only by seahaven-slack-bot) holds. Add moto-backed tests covering un-cancel suppression with field enrichment, status-less merge onto a Cancelled PO, legitimate status updates on non-cancelled POs, new_po backfill of a Cancelled skeleton, fresh create/merge, and authoritative save_cancellation. Refs: #97
2026-07-15 20:17:46 -04:00
"""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)
Merge workorder-ingest into unified procurement repo (#22) * Merge workorder-ingest pipeline into unified repo Move PO lambdas under lambdas/po/, add WO pipeline under lambdas/wo/. Two independent CloudFormation stacks in one CDK app. Fix WO stack compliance: ARM64 architecture, 60-day log retention, aarch64 bundling, RETAIN on Anthropic secret. Remove stale CodePipeline buildspec. * Fix test_local.py import path and remove dead shared/models.py test_local.py referenced the old lambdas/email_processor path. Updated to lambdas/wo/email_processor. Removed shared/ directory entirely as nothing imports from it. * Escape HTML in both web UI dashboards to prevent XSS Both Function URLs are public (auth_type=NONE) and render email-derived content via f-strings. Attacker-crafted emails could inject scripts. Added html.escape() on all interpolated values in both PO and WO dashboards. * Add pagination to WO web UI scan get_work_orders() only fetched the first 1MB page from DynamoDB. Loop on LastEvaluatedKey to match the PO web UI pattern. * Fix esc(None) TypeError and javascript: scheme in PO web UI Coerce supplier name through `or ""` before escaping to handle nested None from DynamoDB. Add scheme allowlist on view_order_url to block javascript:/data: hrefs from LLM-extracted URLs. * Fix WO render_badge None guard, updated_at slice, and backfill path Add null guard to WO render_badge matching the PO version. Use `or ""` before slicing updated_at to handle explicit None values. Fix backfill_sites.py sys.path to use new lambdas/po/site_extractor. * Harden WO web UI and fix JS-context XSS in both dashboards - Use json.dumps for onclick URLs to prevent JS string breakout - Add .lower() to WO render_badge color lookup matching PO pattern - Add pagination to get_comments query - Cap get_work_orders to 500 results matching PO pattern * Apply ruff formatting to web UI handlers
2026-05-12 15:21:06 -04:00
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"}