mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 10:43:14 +00:00
144 lines
6.6 KiB
Python
144 lines
6.6 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.
|
||
|
||
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"}
|