procurement-ingest/lambdas/wo/email_processor/handler.py
Adam Moussa f5c09757f3
feat(lambda): report unhandled Lambda errors to Sentry (PLAT-138) (#203)
* feat(lambda): report unhandled Lambda errors to Sentry

* fix(lambda): strip Sentry stack-frame locals

* style(tests): wrap long lines in sentry_init tests
2026-08-29 20:02:54 +00:00

144 lines
6.6 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
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.
Phase 5: this handler is the thin event loop + fail-closed auth + routing. The
work has moved to flat sibling modules (bare-name imports resolve via the same
flat-landing bundling as ses_auth/template_parser):
extraction.py -- extract_with_bedrock + the <email>-tag neutralizer
telemetry.py -- emit_parse_metric (stdout EMF)
persistence.py -- save_work_order / save_event / comment_id determinism
prompts.py -- EXTRACTION_PROMPT (re-exported below for tests)
"""
import logging
import re
import boto3
import sentry_init # noqa: F401
from email_parsing import parse_raw_email
from extraction import extract_with_bedrock
from persistence import save_event, save_work_order
from prompts import EXTRACTION_PROMPT # noqa: F401 (re-export for tests)
from ses_auth import authenticate_inbound_email
from telemetry import emit_parse_metric
from template_parser import try_deterministic_parse, validate_ai_fallback
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# Lazy cached S3 client. Keeps the public attribute name ``s3`` so the test
# monkeypatch target changes module only, not attribute name.
s3 = None
def _get_s3():
global s3
if s3 is None:
s3 = boto3.client("s3")
return s3
def handler(event, context):
"""Lambda entry point. Triggered by S3 ObjectCreated events."""
# Deploy-guard healthcheck (Phase 0): a top-level direct-invoke
# {"healthcheck": true} probe returns immediately, BEFORE any S3 fetch,
# SES sender-auth gate, or Records iteration. Real mail arrives as S3
# ObjectCreated events whose top-level keys ("Records") AWS controls, so
# email content can never set this key -- this creates no accept path for
# mail. It emits NO EMF and no log line matching the sender_auth_rejected
# metric-filter, so repeated post-deploy smoke invokes never page.
if isinstance(event, dict) and event.get("healthcheck") is True:
return {"healthcheck": "ok"}
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 = _get_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:
method = "ai_fallback"
# A Bedrock transport error (throttle, malformed response, etc.)
# previously emitted ZERO ParseOutcome datapoints -- the only emit
# sites are the post-gate success (below) and the ai_fallback_rejected
# branch. Wrap the call so a fallback attempt that dies in Bedrock
# records exactly one datapoint (ai_fallback / ReasonCode=bedrock_error)
# and then re-raises into the async-retry / DLQ path. The emit is in
# the except -- never pre-call -- so a gate-rejected email (Bedrock
# returned, validate_ai_fallback fails below) still emits ONLY
# ai_fallback_rejected, preserving the wo_stack "a rejected email
# emits nothing else" alarm contract (no double-count).
try:
parsed = extract_with_bedrock(email_data)
except Exception:
emit_parse_metric("ai_fallback", template_id, "bedrock_error", None)
raise
# Fail-closed validation gate on AI output: a prompt-injected
# email body could steer the model into returning arbitrary
# field values, so enforce the same structural contract on both
# parse paths BEFORE any DynamoDB write.
ok, val_reason = validate_ai_fallback(parsed)
if not ok:
logger.warning(
f"AI-fallback validation failed ({val_reason}), skipping: {key}"
)
emit_parse_metric(
"ai_fallback_rejected",
template_id,
val_reason,
parsed.get("work_order_id") if isinstance(parsed, dict) else None,
)
continue
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.
# [0-9] not \d: \d is Unicode-aware and would admit fullwidth digits
# (e.g. "12345") as a distinct-but-lookalike partition key.
work_order_id = parsed.get("work_order_id")
if not work_order_id or not re.fullmatch(r"[0-9]+", 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"}