mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 13:03:14 +00:00
Some checks are pending
Deploy / deploy (push) Waiting to run
Both email-processor God-handlers split along the seams that already work in the flat-sibling pattern established by lambdas/shared/, so bare-name imports keep working under the existing bundling glob. PO (5-way split): handler.py keeps only the event loop, fail-closed auth, and email_type routing. extraction.py holds extract_with_claude and _EMAIL_TAG_RE, importing EXTRACTION_PROMPT from prompts.py and parse_raw_email from shared/email_parsing.py rather than recreating a PO-local copy. enrichment.py is a pure code move of enrich_parsed and pad_zip (PO-only; WO has no enrichment stage) with zero behavior change. telemetry.py holds the EMF ParseMethod emit wrappers. persistence.py holds _write_fields/_merge_update/save_*, collapsing the byte-identical save_new_po/save_revision bodies into one _save_merge helper that both now call through, preserving the sticky Cancelled ConditionExpression guard for both callers; save_cancellation stays separate. WO (5 concerns, no enrichment stage): the handler loop keeps validate_ai_fallback and the re.fullmatch(r"[0-9]+", work_order_id) key guard ahead of both save_work_order and save_event, since the guard protects the DynamoDB partition key and the '#'-delimited comment_id range-key segment. _header_date_iso and comment_id determinism stay colocated with persistence.py's save_event for the retry-idempotent event_id key. EXTRACTION_PROMPT (PO) moves to prompts.py with cross-reference headers to derived_fields.py's authoritative trade/site/fiscal rule tables; handler.py re-exports it (from prompts import EXTRACTION_PROMPT) since four tests dereference handler.EXTRACTION_ PROMPT directly. WO's prompt moves the same way. I/O modules (extraction.py's bedrock client, persistence.py's dynamodb resource, handler.py's s3 client) get lazy cached boto3 accessors; pure modules (enrichment.py, prompts.py, telemetry.py) import no boto3. Test monkeypatch surfaces move to the module that now owns the client (e.g. persistence.dynamodb) everywhere tests patch it, and the moto-before-handler-import ordering in _po_parser_support.py is preserved so the moto-backed suites don't hit real AWS. Behavior-preservation pins, verified with tests: PO still emits ParseMethod=ai_fallback before the Bedrock call, with ai_fallback_rejected as the additive second datapoint on rejection. WO still emits after its gate with mutually-exclusive ai_fallback / ai_fallback_rejected. Shadow DerivedFieldAgreement telemetry stays ai_fallback-only. derived_fields.py is untouched (diff against feature/phase-3-shared-extraction is empty). handler(event, context) signatures and the save_* public contract are unchanged on both pipelines; goldens unchanged. PO_EXPECTED_TOP_LEVEL_MODULES and its WO equivalent in tests/test_bundle_consistency.py are updated for the new sibling modules so the AST bundle-consistency test still fails on an unshipped or uncommented-out sibling.
147 lines
6.1 KiB
Python
147 lines
6.1 KiB
Python
"""
|
|
PO email processor Lambda.
|
|
|
|
Triggered by S3 events when SES delivers a Coupa PO email.
|
|
Parses the raw email, tries the deterministic template parser first, falls
|
|
back to Claude on Bedrock for structured extraction on a miss/invalid result,
|
|
then writes the result to the purchase-orders DynamoDB table.
|
|
|
|
The concerns are split across flat sibling modules (all bundled into the same
|
|
Lambda asset, so bare-name imports resolve):
|
|
* extraction.py -- extract_with_claude + the Bedrock client
|
|
* enrichment.py -- enrich_parsed + pad_zip + derived-field classification
|
|
* telemetry.py -- the EMF ParseMethod / derived-agreement emit wrappers
|
|
* persistence.py -- the purchase-orders merge-writes (save_*)
|
|
* prompts.py -- EXTRACTION_PROMPT (re-exported below for tests)
|
|
This module keeps the S3 event loop, fail-closed sender auth, and email_type
|
|
routing.
|
|
"""
|
|
|
|
import logging
|
|
|
|
import boto3
|
|
from email_parsing import parse_raw_email
|
|
from enrichment import enrich_parsed
|
|
from extraction import extract_with_claude
|
|
from persistence import save_cancellation, save_new_po, save_revision
|
|
|
|
# EXTRACTION_PROMPT is re-exported so handler.EXTRACTION_PROMPT still resolves
|
|
# for the tests that dereference it as a module attribute.
|
|
from prompts import EXTRACTION_PROMPT # noqa: F401
|
|
from ses_auth import authenticate_inbound_email
|
|
from telemetry import _emit_parse_method_metric
|
|
from template_parser import try_deterministic_parse, validate_ai_fallback
|
|
|
|
logger = logging.getLogger()
|
|
logger.setLevel(logging.INFO)
|
|
|
|
# Lazily-built, cached S3 client. Public name ``s3`` is preserved so the
|
|
# monkeypatch attribute is unchanged; building at first CALL (not import) keeps
|
|
# the moto-before-handler invariant and honors any patched fake.
|
|
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."""
|
|
# Direct-invoke healthcheck (post-deploy smoke). This MUST be the very first
|
|
# thing handler() does -- before any S3 fetch, before ses_auth, before the
|
|
# Records loop -- so it (a) creates no accept path for mail (real mail is an
|
|
# S3 ObjectCreated event whose top-level keys AWS controls; email content
|
|
# can never set a top-level "healthcheck" key), and (b) emits no EMF metric
|
|
# and no log line that could match the sender_auth_rejected metric-filter
|
|
# pattern, so two deploys in ~30 min never page that alarm.
|
|
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}")
|
|
|
|
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 purchase 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
|
|
|
|
email_data = parse_raw_email(raw_email)
|
|
logger.info(f"Subject: {email_data['subject']}")
|
|
|
|
# Deterministic template parse first; fall back to the Bedrock AI
|
|
# extractor only on a miss or an invalid (fail-closed) result. The
|
|
# metric is emitted before the Bedrock call, so a Bedrock-side error
|
|
# still records the ai_fallback outcome (the invocation then errors
|
|
# into the errors alarm / DLQ as before).
|
|
parsed, parse_method, template_id, parse_reason = try_deterministic_parse(
|
|
email_data
|
|
)
|
|
_emit_parse_method_metric(
|
|
parse_method,
|
|
template_id,
|
|
parse_reason,
|
|
parsed.get("po_number") if parsed else None,
|
|
)
|
|
if parsed is None:
|
|
parsed = extract_with_claude(email_data)
|
|
# Fail-closed gate on the raw model output (#104 parity): runs
|
|
# BEFORE enrich_parsed and BEFORE any dispatch/save. INTENTIONAL
|
|
# DOUBLE-COUNT: ParseMethod=ai_fallback was already emitted above,
|
|
# BEFORE the Bedrock call (deliberate -- a Bedrock-side error must
|
|
# still record the outcome), so a rejected email produces BOTH an
|
|
# ai_fallback and an ai_fallback_rejected datapoint. The po_stack
|
|
# fallback-rate alarm therefore EXCLUDES the rejected series from
|
|
# its rate math (fb already counts these emails once); see
|
|
# cdk/po_stack.py and README.
|
|
ok, val_reason, normalized = validate_ai_fallback(parsed)
|
|
if not ok:
|
|
logger.warning(
|
|
f"AI-fallback output rejected by validation gate "
|
|
f"({val_reason}); skipping: {key}"
|
|
)
|
|
_emit_parse_method_metric(
|
|
"ai_fallback_rejected",
|
|
template_id,
|
|
val_reason,
|
|
str(parsed.get("po_number"))[:64]
|
|
if isinstance(parsed, dict) and parsed.get("po_number")
|
|
else None,
|
|
)
|
|
# Skip, never raise: attacker-controlled input must not churn
|
|
# the retry/DLQ path.
|
|
continue
|
|
parsed = normalized
|
|
logger.info(
|
|
f"Parsed ({parse_method}/{template_id}/{parse_reason}): "
|
|
f"type={parsed.get('email_type')}, po={parsed.get('po_number')}"
|
|
)
|
|
|
|
if not parsed.get("po_number"):
|
|
logger.warning(f"No PO number found in email, skipping: {key}")
|
|
continue
|
|
|
|
parsed = enrich_parsed(
|
|
parsed, s3_key, email_data["subject"], parse_method=parse_method
|
|
)
|
|
|
|
email_type = parsed.get("email_type")
|
|
if email_type == "cancellation":
|
|
save_cancellation(parsed)
|
|
elif email_type == "revision":
|
|
save_revision(parsed)
|
|
else:
|
|
save_new_po(parsed)
|
|
|
|
return {"statusCode": 200, "body": "OK"}
|