procurement-ingest/lambdas/po/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

179 lines
7.2 KiB
Python

"""PO persistence: merge-writes to the purchase-orders DynamoDB table.
Owns the SET-only merge path (_write_fields / _merge_update) with the sticky
"Cancelled" ConditionExpression guard, the collapsed _save_merge behind the
save_new_po/save_revision public wrappers, and save_cancellation. Uses a lazily
built, cached DynamoDB resource kept under the public name ``dynamodb`` so the
tests' setattr(persistence, "dynamodb", fake) patch surface is unchanged.
"""
import logging
import os
from datetime import datetime, timezone
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders")
# "Cancelled" is a sticky, authoritative status: once a PO reaches it, a later
# new_po/revision may enrich other fields but must never move it back to a
# non-cancelled status.
CANCELLED_STATUS = "Cancelled"
# Lazily-built, cached DynamoDB resource. Public name ``dynamodb`` 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.
dynamodb = None
def _get_dynamodb():
global dynamodb
if dynamodb is None:
dynamodb = boto3.resource("dynamodb")
return dynamodb
def _write_fields(po_number: str, fields: dict, *, guard_cancelled: bool):
"""SET the given non-null fields on a PO record via update_item.
Only the fields supplied are written; absent fields are left untouched, so a
partial payload can never delete data that an earlier email established. The
record is created if it does not exist (DynamoDB update_item upsert).
When ``guard_cancelled`` is True the write carries a ConditionExpression that
only permits it while the record is not already Cancelled. The condition is
evaluated atomically by DynamoDB at write time, so a cancellation that lands
first always wins — there is no read-then-write TOCTOU window. A failed guard
raises ConditionalCheckFailedException for the caller to handle.
"""
table = _get_dynamodb().Table(PO_TABLE)
set_parts = []
attr_names = {}
attr_values = {}
for key, value in fields.items():
if value is None or key == "po_number":
continue
name_ph = f"#{key}"
val_ph = f":{key}"
attr_names[name_ph] = key
attr_values[val_ph] = value
set_parts.append(f"{name_ph} = {val_ph}")
if not set_parts:
return
params = {
"Key": {"po_number": po_number},
"UpdateExpression": "SET " + ", ".join(set_parts),
"ExpressionAttributeNames": attr_names,
"ExpressionAttributeValues": attr_values,
}
if guard_cancelled:
params["ExpressionAttributeValues"][":__cancelled_marker"] = CANCELLED_STATUS
params["ConditionExpression"] = (
"attribute_not_exists(po_status) OR po_status <> :__cancelled_marker"
)
table.update_item(**params)
def _merge_update(po_number: str, fields: dict):
"""Merge (SET-only) the given fields onto a PO record, keeping Cancelled sticky.
Only the fields supplied are written; absent fields are left untouched. The
record is created if it does not exist (DynamoDB update_item upsert).
"Cancelled" is a sticky, authoritative status. When the incoming payload
carries a non-cancelled ``po_status``, the write is guarded by a
ConditionExpression so the status is only applied while the record is not
already Cancelled — enforced atomically at write time, eliminating the
read-then-write TOCTOU where a concurrently-landing cancellation could be
silently un-cancelled. If the guard fails (the PO is already Cancelled), the
same fields are re-written WITHOUT po_status/cancelled_at and
unconditionally, so the other fields still merge while the Cancelled status
stays intact.
A payload with no ``po_status``, or one whose status is already "Cancelled",
needs no guard — a plain merge is correct. This is what keeps legitimate
status updates (non-cancelled PO) and status-less revisions from ever being
dropped: the status is only ever suppressed on a true un-cancel transition.
"""
incoming_status = fields.get("po_status")
if incoming_status is None or incoming_status == CANCELLED_STATUS:
_write_fields(po_number, fields, guard_cancelled=False)
return
try:
_write_fields(po_number, fields, guard_cancelled=True)
except _get_dynamodb().meta.client.exceptions.ConditionalCheckFailedException:
logger.info(
f"PO {po_number} is Cancelled; suppressing incoming "
f"po_status={incoming_status!r} and merging remaining fields"
)
enrich_fields = {
k: v for k, v in fields.items() if k not in ("po_status", "cancelled_at")
}
_write_fields(po_number, enrich_fields, guard_cancelled=False)
def _save_merge(parsed: dict, log_verb: str):
po_number = parsed["po_number"]
fields = {k: v for k, v in parsed.items() if v is not None}
_merge_update(po_number, fields)
logger.info(f"{log_verb} PO {po_number}")
def save_new_po(parsed: dict):
"""Create a PO, merging into any pre-existing record.
Uses a merge update rather than a conditional put so that an out-of-order
cancellation (which leaves a Cancelled skeleton) is filled in with the full
PO data instead of the new_po being silently dropped. "Cancelled" is a sticky
status enforced atomically inside _merge_update: if the PO was already
cancelled, the new_po backfills its remaining fields (supplier, line_items,
amounts) but never un-cancels it.
"""
_save_merge(parsed, "Created/merged")
def save_revision(parsed: dict):
"""Merge revised data into an existing PO without deleting omitted fields.
A revision email often omits unchanged sections (line_items, supplier). The
previous full-overwrite put_item permanently dropped those. This SETs only the
fields present in the revision, leaving everything else intact.
"Cancelled" is a sticky status: a revision may enrich a cancelled PO's fields
but must never move it to a non-cancelled status. That invariant is enforced
atomically inside _merge_update and applies ONLY to the un-cancel transition —
a revision that carries no status change, or one targeting a non-cancelled PO,
updates po_status normally.
"""
_save_merge(parsed, "Revised")
def save_cancellation(parsed: dict):
"""Mark a PO Cancelled, creating a minimal skeleton if it doesn't exist yet.
If the cancellation arrives before the new_po, the skeleton it creates is
later backfilled by save_new_po (which preserves this Cancelled status), so no
PO data is lost on out-of-order delivery.
"""
table = _get_dynamodb().Table(PO_TABLE)
table.update_item(
Key={"po_number": parsed["po_number"]},
UpdateExpression="SET po_status = :status, cancelled_at = :cancelled_at, raw_s3_key = :s3_key",
ExpressionAttributeValues={
":status": "Cancelled",
":cancelled_at": parsed.get(
"processed_at", datetime.now(timezone.utc).isoformat()
),
":s3_key": parsed.get("raw_s3_key", ""),
},
)
logger.info(f"Cancelled PO {parsed['po_number']}")