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

96 lines
4.3 KiB
Python

"""PO parse-outcome and derived-field-agreement EMF telemetry.
Pure module: the emit wrappers below print CloudWatch EMF log lines via the
shared ``emf`` writers (stdout only, no PutMetricData API call), so this module
makes no AWS call and imports no boto3. The derived-agreement wrappers live here
(not in the untouchable derived_fields.py) and are imported by enrichment.py.
"""
import logging
from emf import emit_metric, emit_parse_outcome
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# CloudWatch EMF namespace/metric for the parse-outcome metric (the PO
# fallback-rate alarm in cdk/po_stack.py reads the ["ParseMethod"] series).
METRIC_NAMESPACE = "Seahaven/PoIngest"
# Shadow telemetry for the Python-derived classifier bake. One EMF record per
# derived field per email, emitted ONLY on the ai_fallback path (the template
# path has no LLM value to compare against). Dimensioned by Field x Agreement
# only -- PythonValue/LlmValue/po_number ride along as Logs-Insights-queryable
# properties so the cardinality stays fixed at (3 fields x 4 categories).
DERIVED_METRIC_NAME = "DerivedFieldAgreement"
def _emit_parse_method_metric(method, template_id, reason_code, po_number):
"""Emit one CloudWatch EMF line recording the parse outcome.
Zero-latency (no PutMetricData API call): the extraction path is async and
the role already has logs:PutLogEvents. ParseMethod/TemplateId are the only
promoted (dimensioned) fields to keep cardinality low; ReasonCode and
po_number ride along as Logs-Insights-queryable properties.
Two dimension sets are published: ["ParseMethod"] (aggregated across all
template ids -- the series the fallback-rate alarm queries) AND
["ParseMethod", "TemplateId"] (per-template breakdown for Logs Insights /
dashboards). CloudWatch materializes only the exact dimension sets listed
here and does NOT auto-aggregate, so the alarm's single-dimension query
would receive no data unless ["ParseMethod"] is emitted explicitly."""
emit_parse_outcome(
METRIC_NAMESPACE, method, template_id, reason_code, "po_number", po_number
)
def _derived_agreement(llm_value, python_value) -> str | None:
"""Classify Python-vs-LLM agreement for one derived field.
Returns None when both values are None (nothing to compare -- the caller
then skips emission). Categories:
* ``agree`` -- both non-None and equal after str-strip
* ``disagree`` -- both non-None but different
* ``llm_null_python_filled``-- LLM None, Python supplied a value
* ``python_null`` -- LLM non-None, Python None
"""
if llm_value is None and python_value is None:
return None
if llm_value is None:
return "llm_null_python_filled"
if python_value is None:
return "python_null"
if str(llm_value).strip() == str(python_value).strip():
return "agree"
return "disagree"
def _emit_derived_agreement_metric(field, llm_value, python_value, po_number):
"""Emit one CloudWatch EMF line shadowing the Python-derived classifier
against the LLM value for a single derived field (ai_fallback path only).
Mirrors ``_emit_parse_method_metric``: zero-latency (no PutMetricData; the
role already has logs:PutLogEvents), Field x Agreement the only promoted
dimension set (cardinality 3x4). PythonValue/LlmValue/po_number ride along
as Logs-Insights-queryable properties so a disagreement can be reviewed by
example without inflating metric cardinality. No emission when both values
are None -- there is nothing to compare."""
agreement = _derived_agreement(llm_value, python_value)
if agreement is None:
return
emit_metric(
METRIC_NAMESPACE,
DERIVED_METRIC_NAME,
[["Field", "Agreement"]],
{
"Field": field,
"Agreement": agreement,
"po_number": po_number or "",
# Length-clamped: Python values are regex/enum-bounded by construction,
# but the LLM value is schema-unvalidated model output -- a hallucinated
# free-text field must not land unbounded in a 2-month log line
# (sh-security-review PO-DC-02, confirmed low).
"PythonValue": "" if python_value is None else str(python_value)[:64],
"LlmValue": "" if llm_value is None else str(llm_value)[:64],
},
)