procurement-ingest/lambdas/wo/email_processor/persistence.py
Adam Moussa ca4f43a2cc
Some checks are pending
Deploy / deploy (push) Waiting to run
feat: decompose email-processor handlers into flat siblings + lazy boto3 clients (refactor phase 5) (#113)
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.
2026-07-20 15:34:53 -04:00

173 lines
6.4 KiB
Python

"""DynamoDB persistence for the work-order email processor.
Owns the work-orders + comments table writes and the retry-idempotent event_id /
comment_id determinism (_header_date_iso stays coupled with save_event so the
key derivation is single-sourced). The boto3 DynamoDB resource is built lazily on
first use so tests can patch this module's ``dynamodb`` attribute before any real
client is constructed (moto-before-handler invariant).
"""
import hashlib
import logging
import os
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
# Lazy cached DynamoDB resource. Keeps the public attribute name ``dynamodb`` so
# the test monkeypatch target changes module only, not attribute name.
dynamodb = None
def _get_dynamodb():
global dynamodb
if dynamodb is None:
dynamodb = boto3.resource("dynamodb")
return dynamodb
def save_work_order(parsed: dict, s3_key: str):
"""Create or update a work order in DynamoDB."""
table = _get_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 = _get_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})")