mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 08:23:14 +00:00
* Add deterministic template parser for WO emails The workorder-email-processor sends every one of ~22.9k emails/month to an LLM, but ~93.6% are the plain-text "AMAZON UPDATE WO DETAILS" comment template and ~6.4% the HTML "AMAZON assign Work Order" template. Parse those two shapes deterministically, offline, so the AI call is reserved for the long tail. The module is pure (no boto3, no network). try_deterministic_parse classifies by subject, extracts the shared contract fields, and returns a result ONLY when it passes a strict fail-closed validation gate: exact contract-key set, subject/id agreement, the literal "Work Order: <id>" double space, per-type required fields, site-code shape, and a label-bleed guard so a value that over-ran into the next field fails. Any miss, drift, or extractor exception yields None so the caller falls back to the AI extractor -- data is never corrupted, only the fallback rate rises. Refs: #23 * Migrate WO processor to Bedrock and fix comment_id collision Switch the AI path from the Anthropic SDK to bedrock-runtime InvokeModel on the inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0 (BEDROCK_MODEL_ID env), so parsing no longer needs a provider API key or Secrets Manager secret. The EXTRACTION_PROMPT and JSON contract are kept byte-identical, so the AI-fallback output is unchanged. Try the new deterministic template parser first and only call Bedrock on a miss/invalid result. Fix issue #23: the WorkOrderComments range key was work_order_id#<comment_time>, so two emails on one WO with an identical or absent comment time collided and overwrote each other. Derive a 12-hex suffix from the S3 object key alone -- deterministic, so an async retry of the same object is byte-identical (idempotent) while distinct emails get distinct keys -- and keep wall-clock now() out of the key (literal 'nocomment' segment when comment_time is absent). Also emit one CloudWatch EMF line per record (Seahaven/WorkorderIngest ParseOutcome, dimensioned by ParseMethod/TemplateId) for parse-outcome observability, replace the deprecated datetime.utcnow() with datetime.now(timezone.utc), and drop the anthropic dependency. Refs: #23 * Migrate PO processor to Bedrock Switch the PO email processor's AI extraction from the Anthropic SDK to bedrock-runtime InvokeModel on the inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0 (BEDROCK_MODEL_ID env), so it no longer needs a provider API key or Secrets Manager secret. PO parsing stays fully AI -- only the provider changes. The EXTRACTION_PROMPT is kept byte-identical and the Bedrock text output is still decoded with json.loads(..., parse_float=Decimal), which DynamoDB requires (it rejects floats). Replace the deprecated datetime.utcnow() with datetime.now(timezone.utc) and drop the anthropic dependency. * Grant Bedrock IAM, drop Anthropic secrets, add fallback alarm Both stacks moved their processors from the Anthropic API to the Bedrock inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0. Grant each processor role bedrock:InvokeModel + bedrock:InvokeModelWithResponseStream on BOTH the inference-profile ARN AND the per-region foundation-model ARNs for us-east-1/us-east-2/us-west-2 (empty-account) -- the us.* profile routes cross-region, so a profile-only grant AccessDenies at runtime. Remove both anthropic-api-key Secret constructs, their grant_read, and the ANTHROPIC_API_KEY_SECRET_ARN env; add BEDROCK_MODEL_ID. The secrets had RemovalPolicy.RETAIN so they are orphaned, not deleted -- flagged in the README for manual post-deploy deletion and key revocation. Add the workorder-email-processor-template-fallback-rate alarm: a FILL(0) + >=10-sample volume-floor MathExpression over the EMF ParseOutcome metric (15-min periods) that pages when the AI-fallback share exceeds 15% sustained, catching Hexagon template drift. ALARM-only SnsAction to site-alerts, no OK action, NOT_BREACHING, matching the existing stack idiom. * Add offline WO parser test suite Cover the deterministic parser with golden-file tests over 55 real scrubbed .eml fixtures (both comment sub-shapes, username Submitted-By, address present/absent, br+CRLF assign addresses), fail-closed validation-gate rules, adversarial and prompt-injection cases that must route to ai_fallback or parse without corrupting other fields, the issue #23 comment_id idempotency invariants, and the Bedrock-fallback dispatch plus EMF-metric emission with a mocked invoke_model. Extend pytest.ini testpaths to discover the co-located suite, and update tests/conftest.load_handler to put a handler's own directory on sys.path so the WO handler's new `from template_parser import ...` resolves under the existing shared handler tests. Point test_local.py at the new template-first + Bedrock flow. Refs: #23 * Document Bedrock migration and WO parse flow in README Record the provider switch to the Bedrock inference profile (no Anthropic API key or Secrets Manager secret, with the retired secrets flagged for manual deletion), the WO deterministic-template-first + AI-fallback flow, the new ParseOutcome EMF metric and template-fallback-rate alarm, the issue #23 comment_id format change, the +00:00 aware-UTC timestamp shift, and offline test instructions. Refs: #23 * Fix f-string lint and formatting in backfill scripts Drop the f prefix from two f-strings that carry no placeholders (F541) and apply ruff format, so `ruff check` / `ruff format --check` pass in CI. * Emit ParseMethod-only EMF set so fallback alarm can fire The fallback-rate alarm queries the ParseOutcome series keyed on ParseMethod alone, but the emitter published only the joint (ParseMethod, TemplateId) dimension set. CloudWatch materializes exactly the listed dimension sets and does not auto-aggregate, so the alarm's series never received data: it evaluated a constant 0 and could never page on template-drift coverage collapse. Publish both ["ParseMethod"] and ["ParseMethod","TemplateId"] and update the EMF regression test to assert both sets are present. * Commit WO parser .eml fixtures for executable coverage The parser test suite globbed for input .eml fixtures that the repo's `*.eml` ignore rule kept uncommitted, so every parametrized golden and fail-closed test collected zero cases and CI could not exercise the deterministic parser that handles 100% of WO email volume. Add a fixtures-only negation to .gitignore and commit the 55 scrubbed positive samples (50 update-plaintext, 5 assign-html) plus 14 ai-fallback and 3 adversarial fixtures. The ai-fallback set covers each fail-closed reason code (subject_no_match, single_space_work_order, malformed_site_code, label_bleed, creation_time_unparseable, wo_id_mismatch, missing_required_field) and the adversarial set proves the parser is total and confines prompt-injection payloads to comment_text without steering the structured fields. * Fix WO parser advisories A1-A3 (PR #99 follow-ups) A1 — AI-fallback comment_id nondeterminism: parsed comment_time is model output and not stable across Lambda async retries, so on the ai_fallback path the comment_id range-key time segment now derives from the email Date header (deterministic per S3 object) instead of the model's comment_time. The template path is unchanged (its comment_time is a pure function of the raw email). Bedrock invoke pins temperature 0 so retries reproduce the same extraction. Closes the #23 reopening on the AI path. A2 — EMF record now carries the spec-required _aws.Timestamp (epoch ms) so CloudWatch reliably extracts the ParseOutcome datapoint that the fallback-rate alarm depends on. A3 — T1 New Comment capture no longer truncates at the first blank line; multi-paragraph comments are captured through internal blanks and terminate at the next label/separator. 17 golden files regenerated from the real fixtures accordingly. Hardening from the sh-security-review pass on this diff: - _header_date_iso is total: OverflowError/OSError from an extreme Date header fall back to 'nocomment' instead of failing the invocation. - _capture_block trims blanks in O(n) (no pop(0)) — removes a quadratic path on a crafted large blank run. - work_order_id is enforced digits-only on BOTH parse paths before it is used as a DynamoDB key, so prompt-injected AI output cannot forge '#' range-key segments or land on an arbitrary WO.
385 lines
15 KiB
Python
385 lines
15 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 hashlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import re
|
|
from datetime import datetime, timezone
|
|
from email import policy
|
|
from email.utils import parsedate_to_datetime
|
|
|
|
import boto3
|
|
from ses_auth import authenticate_inbound_email
|
|
|
|
from template_parser import try_deterministic_parse
|
|
|
|
logger = logging.getLogger()
|
|
logger.setLevel(logging.INFO)
|
|
|
|
s3 = boto3.client("s3")
|
|
dynamodb = boto3.resource("dynamodb")
|
|
bedrock = boto3.client("bedrock-runtime")
|
|
|
|
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
|
|
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
|
|
BEDROCK_MODEL_ID = os.environ.get(
|
|
"BEDROCK_MODEL_ID", "us.anthropic.claude-haiku-4-5-20251001-v1:0"
|
|
)
|
|
|
|
# CloudWatch EMF namespace + metric for parse-outcome observability.
|
|
METRIC_NAMESPACE = "Seahaven/WorkorderIngest"
|
|
METRIC_NAME = "ParseOutcome"
|
|
|
|
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 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_bedrock(email_data: dict) -> dict:
|
|
"""Send parsed email to Claude on Bedrock for structured extraction."""
|
|
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']}"
|
|
)
|
|
|
|
resp = bedrock.invoke_model(
|
|
modelId=BEDROCK_MODEL_ID,
|
|
body=json.dumps(
|
|
{
|
|
"anthropic_version": "bedrock-2023-05-31",
|
|
"max_tokens": 1024,
|
|
# Greedy decoding: retries of the same email should get the
|
|
# same extraction back (advisory A1). Not a hard guarantee of
|
|
# determinism, so model output still never enters a table key.
|
|
"temperature": 0,
|
|
"messages": [
|
|
{
|
|
"role": "user",
|
|
"content": f"{EXTRACTION_PROMPT}\n\nEMAIL:\n{email_text}",
|
|
}
|
|
],
|
|
}
|
|
),
|
|
)
|
|
response_text = json.loads(resp["body"].read())["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 emit_parse_metric(method, template_id, reason_code, work_order_id):
|
|
"""Emit one CloudWatch EMF line recording the parse outcome.
|
|
|
|
Zero-latency (no PutMetricData API call): the extraction path is async and
|
|
the role already has logs:PutLogEvents. ParseMethod/TemplateId are the only
|
|
promoted (dimensioned) fields to keep cardinality low; ReasonCode and
|
|
work_order_id ride along as Logs-Insights-queryable properties.
|
|
|
|
Two dimension sets are published: ["ParseMethod"] (aggregated across all
|
|
template ids -- the series the fallback-rate alarm queries) AND
|
|
["ParseMethod", "TemplateId"] (per-template breakdown for Logs Insights /
|
|
dashboards). CloudWatch materializes only the exact dimension sets listed
|
|
here and does NOT auto-aggregate, so the alarm's single-dimension query
|
|
would receive no data unless ["ParseMethod"] is emitted explicitly."""
|
|
emf = {
|
|
"_aws": {
|
|
# EMF requires Timestamp (epoch ms); without it CloudWatch may not
|
|
# extract the metric datapoint from the log event (advisory A2).
|
|
"Timestamp": int(datetime.now(timezone.utc).timestamp() * 1000),
|
|
"CloudWatchMetrics": [
|
|
{
|
|
"Namespace": METRIC_NAMESPACE,
|
|
"Dimensions": [["ParseMethod"], ["ParseMethod", "TemplateId"]],
|
|
"Metrics": [{"Name": METRIC_NAME, "Unit": "Count"}],
|
|
}
|
|
],
|
|
},
|
|
"ParseMethod": method,
|
|
"TemplateId": template_id or "unknown",
|
|
"ReasonCode": reason_code or "ok",
|
|
"work_order_id": work_order_id or "",
|
|
METRIC_NAME: 1,
|
|
}
|
|
print(json.dumps(emf))
|
|
|
|
|
|
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.now(timezone.utc).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 _header_date_iso(header_date):
|
|
"""Parse an RFC 2822 Date header into a UTC ISO string, or None.
|
|
|
|
Deterministic for a given raw email, unlike model output. Total: any
|
|
unparseable/out-of-range header (incl. OverflowError from extreme years,
|
|
which is NOT a ValueError) yields None, never an exception -- a crafted
|
|
Date header must not be able to fail the invocation."""
|
|
if not header_date:
|
|
return None
|
|
try:
|
|
dt = parsedate_to_datetime(header_date)
|
|
if dt is None:
|
|
return None
|
|
if dt.tzinfo is None:
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
return dt.astimezone(timezone.utc).isoformat()
|
|
except (TypeError, ValueError, OverflowError, OSError):
|
|
return None
|
|
|
|
|
|
def save_event(
|
|
parsed: dict,
|
|
s3_key: str,
|
|
object_key: str,
|
|
parse_method: str,
|
|
header_date: str | None,
|
|
):
|
|
"""Save an event to the events table. Every email creates an event entry.
|
|
|
|
Issue #23: the range key (attr name stays ``comment_id``) must be unique per
|
|
source email AND identical across Lambda async retries of the same S3 object.
|
|
The uniqueness suffix is a deterministic hash of the S3 object key alone (not
|
|
the s3:// URI, so it is stable across a bucket rename), and wall-clock now()
|
|
is kept OUT of the key. When no time is available we use the literal
|
|
'nocomment' segment rather than now() -- otherwise each retry would produce a
|
|
different key and duplicate the row. Two distinct emails on the same WO map
|
|
to distinct object keys -> distinct rows.
|
|
|
|
Advisory A1: the time segment may come from parsed comment_time ONLY on the
|
|
template path, where it is a pure function of the raw email. On the AI path
|
|
the model can return a different comment_time on a retry (even at
|
|
temperature 0 determinism is not guaranteed), which would fork the key and
|
|
duplicate the row -- so there the segment derives from the email's Date
|
|
header instead."""
|
|
table = dynamodb.Table(COMMENTS_TABLE)
|
|
|
|
work_order_id = parsed["work_order_id"]
|
|
email_type = parsed.get("email_type", "unknown")
|
|
key_suffix = hashlib.sha256(object_key.encode("utf-8")).hexdigest()[:12]
|
|
comment_time = parsed.get("comment_time")
|
|
if parse_method == "template":
|
|
time_part = comment_time if comment_time else "nocomment"
|
|
else:
|
|
time_part = _header_date_iso(header_date) or "nocomment"
|
|
event_id = f"{work_order_id}#{time_part}#{key_suffix}"
|
|
# created_at is display-only; may fall back to now() without affecting the key.
|
|
display_time = comment_time or datetime.now(timezone.utc).isoformat()
|
|
|
|
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": display_time,
|
|
"source_email_s3_key": s3_key,
|
|
"ingested_at": datetime.now(timezone.utc).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']}")
|
|
|
|
# Deterministic template parse first; fall back to the AI extractor only
|
|
# on a miss or an invalid (fail-closed) result.
|
|
parsed, method, template_id, reason = try_deterministic_parse(email_data)
|
|
if parsed is None:
|
|
parsed = extract_with_bedrock(email_data)
|
|
method = "ai_fallback"
|
|
logger.info(
|
|
f"Parsed ({method}/{template_id}/{reason}): "
|
|
f"type={parsed.get('email_type')}, wo={parsed.get('work_order_id')}"
|
|
)
|
|
|
|
emit_parse_metric(method, template_id, reason, parsed.get("work_order_id"))
|
|
|
|
# work_order_id becomes a DynamoDB partition key and the leading, '#'-
|
|
# delimited segment of the comment_id range key, so it must be digits
|
|
# only. The template path already guarantees this via validate(); the
|
|
# AI-fallback path returns raw model output, which a prompt-injected
|
|
# email body could steer into a non-numeric or '#'-bearing value that
|
|
# forges key segments or lands on an arbitrary WO. Enforce the same
|
|
# contract on both paths and skip (fail closed) on a violation.
|
|
work_order_id = parsed.get("work_order_id")
|
|
if not work_order_id or not re.fullmatch(r"\d+", str(work_order_id)):
|
|
logger.warning(f"Missing or non-numeric work order ID, 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. The raw object key
|
|
# (not the s3:// URI) drives the retry-idempotent comment_id suffix;
|
|
# the parse method decides whether comment_time may enter the key (A1).
|
|
save_event(parsed, s3_key, key, method, email_data.get("date"))
|
|
|
|
return {"statusCode": 200, "body": "OK"}
|