procurement-ingest/lambdas/wo/email_processor/template_parser.py
Adam Moussa d677358801
Some checks are pending
Deploy / deploy (push) Waiting to run
fix(wo): reject non-str AI free-text fields (#158)
2026-08-04 19:39:20 -04:00

531 lines
20 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Deterministic template parser for Hexagon EAM work-order emails.
Pure module: no boto3, no network. Runs ahead of the AI extraction path in the
work-order email processor. Only returns a parsed result when it is proven
conformant to one of the two known Hexagon templates; otherwise it fails closed
and signals the caller to fall back to the AI extractor.
Two templates (see BUILD SPEC / recon):
T1 update_plaintext -- Subject "AMAZON UPDATE WO DETAILS <id>", text/plain,
labels New Comment: / Creation Time(UTC): / Submitted By: /
"Work Order: <id> - <desc>" (double space) / "Building: <SITE>." + address.
Maps to email_type "comment".
T2 assign_html -- Subject "AMAZON assign Work Order <id> on building <SITE>",
simple HTML, <br>-delimited WO Description: / Severity: / Date Reported: /
Scheduled Start Date: / Address:. Maps to email_type "new_work_order".
Entry point: try_deterministic_parse(email_data) -> (parsed|None, method,
template_id, reason_code).
"""
import html as html_module
import re
# The contract keys (16), EXACTLY -- mirrors the AI EXTRACTION_PROMPT fields. A conformant parse is a dict with these keys
# and no others (validation rule 2).
CONTRACT_KEYS = (
"email_type",
"work_order_id",
"description",
"status",
"site_code",
"building",
"address",
"severity",
"priority",
"date_reported",
"scheduled_start",
"due_date",
"assigned_to",
"commenter",
"comment_text",
"comment_time",
)
VALID_EMAIL_TYPES = {"new_work_order", "update", "comment", "cancellation"}
VALID_STATUSES = {
"new",
"assigned",
"in_progress",
"on_hold",
"completed",
"cancelled",
"unknown",
}
# Free-text scalar fields that must be None or str (blocks LLM-emitted maps/
# lists from landing as DynamoDB Map/List attribute pollution, and floats that
# would crash update_item). email_type, status, work_order_id, site_code, and
# date fields are validated separately.
_AI_FREE_TEXT_STR_FIELDS = (
"description",
"building",
"address",
"severity",
"priority",
"assigned_to",
"commenter",
"comment_text",
)
# Subject classifiers.
_T1_SUBJECT = re.compile(r"^AMAZON UPDATE WO DETAILS\s+(?P<wo>\S+)\s*$")
_T2_SUBJECT = re.compile(
r"^AMAZON assign Work Order\s+(?P<wo>\S+)\s+on building\s+(?P<site>\S+)\s*$"
)
# Known label tokens per template (lower-cased, used for scanning + bleed guard).
_T1_LABELS = (
"new comment:",
"creation time(utc):",
"submitted by:",
"work order:",
"building:",
)
_T2_LABELS = (
"wo description:",
"severity:",
"date reported:",
"scheduled start date:",
"address:",
)
_SEPARATOR_RE = re.compile(r"_{4,}")
# \A...\Z (not ^...$, whose $ also matches just before a trailing newline) and
# explicit [0-9] (not \d, which is Unicode-aware and would admit fullwidth
# digits) so a value like "WIL1\n" or "12345" cannot pass as well-formed.
_SITE_CODE_RE = re.compile(r"\A[A-Z]{2,4}[0-9]{1,2}\Z")
_WO_ID_RE = re.compile(r"\A[0-9]+\Z")
def _empty_candidate():
return {k: None for k in CONTRACT_KEYS}
def _normalize(candidate):
"""Guarantee all contract keys exist (None for absent) before returning."""
out = _empty_candidate()
for k in CONTRACT_KEYS:
if k in candidate:
out[k] = candidate[k]
return out
def classify_template(email_data):
"""Return (template_id, reason). template_id in {update_plaintext,
assign_html, unknown}."""
subject = (email_data.get("subject") or "").strip()
if _T1_SUBJECT.match(subject):
return "update_plaintext", "ok"
if _T2_SUBJECT.match(subject):
return "assign_html", "ok"
return "unknown", "subject_no_match"
def _subject_ids(email_data):
"""Return (wo_id, site_code) parsed from the subject, or (None, None)."""
subject = (email_data.get("subject") or "").strip()
m = _T1_SUBJECT.match(subject)
if m:
return m.group("wo"), None
m = _T2_SUBJECT.match(subject)
if m:
return m.group("wo"), m.group("site")
return None, None
def _matches_any_label(line, labels):
low = line.strip().lower()
return any(low.startswith(lab) for lab in labels)
def _contains_label_or_separator(text, labels):
"""Label-bleed guard: True if text carries a known label token or a
separator run (indicates the value over-ran into the next field)."""
if text is None:
return False
if _SEPARATOR_RE.search(text):
return True
low = text.lower()
return any(lab in low for lab in labels)
def _html_to_lines(body):
"""Tiny HTML->text: turn <br>/</p>/</tr> and real CRLFs into line breaks,
strip remaining tags, unescape entities. Returns a list of raw lines."""
text = re.sub(r"(?i)<br\s*/?>", "\n", body)
text = re.sub(r"(?i)</p\s*>", "\n", text)
text = re.sub(r"(?i)</tr\s*>", "\n", text)
text = re.sub(r"<[^>]+>", "", text)
text = html_module.unescape(text)
text = text.replace("\r\n", "\n").replace("\r", "\n")
return text.split("\n")
def _plain_lines(body):
return body.replace("\r\n", "\n").replace("\r", "\n").split("\n")
def _find_label_index(lines, label):
ll = label.lower()
for i, line in enumerate(lines):
if line.strip().lower().startswith(ll):
return i
return -1
def _inline_value(line, label):
"""Value = remainder of the label line after the label token."""
idx = line.lower().find(label.lower())
return line[idx + len(label) :].strip()
def _capture_block(lines, start_index, labels, stop_on_blank=True):
"""Collect lines after start_index until a blank line (unless
stop_on_blank=False), a separator run, or a known label. Returns a list of
stripped non-consumed lines (may be empty).
stop_on_blank=False is for free-text blocks that legitimately contain blank
lines (multi-paragraph comments, advisory A3): internal blanks are kept as
empty strings, leading/trailing blanks are trimmed."""
collected = []
j = start_index + 1
while j < len(lines):
s = lines[j].strip()
if s == "":
if stop_on_blank:
break
# Leading blanks are never collected, so a long blank run cannot
# accumulate ahead of the O(1)-per-line trim below.
if collected:
collected.append("")
j += 1
continue
if _SEPARATOR_RE.fullmatch(s) or _SEPARATOR_RE.search(s):
break
if _matches_any_label(s, labels):
break
collected.append(s)
j += 1
while collected and collected[-1] == "":
collected.pop()
return collected
# ---------------------------------------------------------------------------
# T1: update_plaintext -> comment
# ---------------------------------------------------------------------------
def extract_update_plaintext(email_data):
"""Extract all contract keys from a T1 plaintext update email. Returns a full
dict (values None where absent). Correctness is enforced by validate()."""
subject_wo, _ = _subject_ids(email_data)
lines = _plain_lines(email_data.get("body") or "")
candidate = _empty_candidate()
candidate["email_type"] = "comment"
candidate["work_order_id"] = subject_wo
# status stays None on a comment upsert -- never clobber a real wo_status.
# --- New Comment: block up to separator / next label. Blank lines do NOT
# end the block (multi-paragraph comments, advisory A3); the next label
# (normally Creation Time(UTC):) is the terminator. ---
nc_idx = _find_label_index(lines, "New Comment:")
if nc_idx >= 0:
pieces = []
inline = _inline_value(lines[nc_idx], "New Comment:")
if inline:
pieces.append(inline)
pieces.extend(_capture_block(lines, nc_idx, _T1_LABELS, stop_on_blank=False))
candidate["comment_text"] = "\n".join(pieces) if pieces else None
# --- Creation Time(UTC): -> ISO ---
ct_idx = _find_label_index(lines, "Creation Time(UTC):")
if ct_idx >= 0:
raw = _inline_value(lines[ct_idx], "Creation Time(UTC):")
candidate["comment_time"] = _parse_dt(raw, "%Y-%m-%d %H:%M:%S")
# --- Submitted By: -> commenter (username or joined ARN) ---
sb_idx = _find_label_index(lines, "Submitted By:")
if sb_idx >= 0:
pieces = []
inline = _inline_value(lines[sb_idx], "Submitted By:")
if inline:
pieces.append(inline)
# ARN continuation lines (rare) join with no separator.
pieces.extend(_capture_block(lines, sb_idx, _T1_LABELS))
candidate["commenter"] = "".join(pieces) if pieces else None
# --- Work Order: <id> - <desc> ---
wo_idx = _find_label_index(lines, "Work Order:")
if wo_idx >= 0:
m = re.match(r"\s*Work Order:\s+(\S+)\s+-\s+(.*)$", lines[wo_idx])
if m:
desc = m.group(2).strip()
candidate["description"] = desc or None
# --- Building: <SITE>. + optional address block ---
b_idx = _find_label_index(lines, "Building:")
if b_idx >= 0:
site = _inline_value(lines[b_idx], "Building:").rstrip(".").strip()
if site:
candidate["site_code"] = site
candidate["building"] = site
addr = _capture_block(lines, b_idx, _T1_LABELS)
candidate["address"] = "\n".join(addr) if addr else None
return _normalize(candidate)
# ---------------------------------------------------------------------------
# T2: assign_html -> new_work_order
# ---------------------------------------------------------------------------
def extract_assign_html(email_data):
"""Extract all contract keys from a T2 HTML assign email."""
subject_wo, subject_site = _subject_ids(email_data)
lines = _html_to_lines(email_data.get("body") or "")
candidate = _empty_candidate()
candidate["email_type"] = "new_work_order"
candidate["status"] = "assigned"
candidate["work_order_id"] = subject_wo
candidate["site_code"] = subject_site
candidate["building"] = subject_site
d_idx = _find_label_index(lines, "WO Description:")
if d_idx >= 0:
desc = _inline_value(lines[d_idx], "WO Description:")
candidate["description"] = desc or None
s_idx = _find_label_index(lines, "Severity:")
if s_idx >= 0:
sev = _inline_value(lines[s_idx], "Severity:")
candidate["severity"] = sev or None
dr_idx = _find_label_index(lines, "Date Reported:")
if dr_idx >= 0:
raw = _inline_value(lines[dr_idx], "Date Reported:")
candidate["date_reported"] = _parse_dt(raw, "%Y-%m-%d %H:%M")
ss_idx = _find_label_index(lines, "Scheduled Start Date:")
if ss_idx >= 0:
raw = _inline_value(lines[ss_idx], "Scheduled Start Date:")
candidate["scheduled_start"] = _parse_date(raw, "%Y-%m-%d")
a_idx = _find_label_index(lines, "Address:")
if a_idx >= 0:
inline = _inline_value(lines[a_idx], "Address:")
pieces = []
if inline:
pieces.append(inline)
pieces.extend(_capture_block(lines, a_idx, _T2_LABELS))
candidate["address"] = "\n".join(pieces) if pieces else None
return _normalize(candidate)
def _parse_dt(raw, fmt):
from datetime import datetime
try:
return datetime.strptime(raw.strip(), fmt).isoformat()
except (ValueError, AttributeError):
return _UNPARSEABLE
def _parse_date(raw, fmt):
from datetime import datetime
try:
return datetime.strptime(raw.strip(), fmt).date().isoformat()
except (ValueError, AttributeError):
return _UNPARSEABLE
# Sentinel: a label was present but its date/time value did not parse. This must
# FAIL the gate (present-but-unparseable), distinct from an absent value (None).
_UNPARSEABLE = "__UNPARSEABLE__"
# ---------------------------------------------------------------------------
# Validation gate -- FAIL CLOSED
# ---------------------------------------------------------------------------
def validate(candidate, template_id, email_data):
"""Return (True, 'ok') only if the candidate is provably conformant; else
(False, reason). Every rule must hold."""
# (1) known template
if template_id not in ("update_plaintext", "assign_html"):
return False, "subject_no_match"
# (2) keys EXACTLY the contract set
if set(candidate.keys()) != set(CONTRACT_KEYS):
return False, "key_set_mismatch"
# Unparseable date sentinels never survive.
for key in ("comment_time", "date_reported", "scheduled_start"):
if candidate.get(key) == _UNPARSEABLE:
return False, "creation_time_unparseable"
subject_wo, subject_site = _subject_ids(email_data)
body = email_data.get("body") or ""
# (3) work_order_id non-empty digits AND == subject id
wo = candidate.get("work_order_id")
if not wo or not _WO_ID_RE.match(str(wo)):
return False, "missing_required_field"
if wo != subject_wo:
return False, "wo_id_mismatch"
# (5) email_type in enum AND == template's expected type
expected_type = "comment" if template_id == "update_plaintext" else "new_work_order"
et = candidate.get("email_type")
if et not in VALID_EMAIL_TYPES:
return False, "missing_required_field"
if et != expected_type:
return False, "email_type_mismatch"
# (6) site_code if set matches the code pattern
site = candidate.get("site_code")
if site is not None and not _SITE_CODE_RE.match(str(site)):
return False, "malformed_site_code"
# (7) status if non-null in the enum
status = candidate.get("status")
if status is not None and status not in VALID_STATUSES:
return False, "invalid_status"
if template_id == "update_plaintext":
# (4) body must contain "Work Order: <id>" with LITERAL double space.
if f"Work Order: {subject_wo}" not in body:
return False, "single_space_work_order"
# (10) if Creation Time present it must have parsed.
if "Creation Time(UTC):" in body and candidate.get("comment_time") is None:
return False, "creation_time_unparseable"
# (8) comment needs work_order_id + non-empty comment_text
if not (candidate.get("comment_text") or "").strip():
return False, "missing_required_field"
# (9) label-bleed guard on comment_text
if _contains_label_or_separator(candidate.get("comment_text"), _T1_LABELS):
return False, "label_bleed"
if _contains_label_or_separator(candidate.get("address"), _T1_LABELS):
return False, "label_bleed"
else: # assign_html
# (4) body id if present must == subject
body_ids = re.findall(r"Work Order\D*(\d+)", body)
for bid in body_ids:
if bid != subject_wo:
return False, "wo_id_mismatch"
# (10) all five labels present and both dates parse.
for label in _T2_LABELS:
if label.lower() not in body.lower():
return False, "missing_required_field"
if candidate.get("date_reported") is None:
return False, "creation_time_unparseable"
if candidate.get("scheduled_start") is None:
return False, "creation_time_unparseable"
# (8) new_work_order needs work_order_id + site_code + non-empty description
if not candidate.get("site_code"):
return False, "missing_required_field"
if not (candidate.get("description") or "").strip():
return False, "missing_required_field"
# (9) label-bleed guard on description + address
if _contains_label_or_separator(candidate.get("description"), _T2_LABELS):
return False, "label_bleed"
if _contains_label_or_separator(candidate.get("address"), _T2_LABELS):
return False, "label_bleed"
return True, "ok"
def validate_ai_fallback(candidate):
"""Fail-closed schema/enum validation for the AI-fallback parse path.
Called on the raw Bedrock/Claude output BEFORE any DynamoDB write. The
contract is a subset of the template-path validate(): no template-specific
rules (subject-id matching, label-bleed, required fields per template),
but ALL structural/enum guards apply so the AI path converges on the same
structural contract as the template path.
Returns (True, 'ok') or (False, reason)."""
from datetime import datetime
# json.loads on model output can yield any JSON type; only an object can
# satisfy the contract, and anything else must fail closed here rather
# than crash the handler into async S3 retries / DLQ.
if not isinstance(candidate, dict):
return False, "not_an_object"
# keys EXACTLY the contract set
if set(candidate.keys()) != set(CONTRACT_KEYS):
return False, "key_set_mismatch"
# work_order_id non-empty digits-only. The handler keeps its own check
# where the DynamoDB key is built (defense in depth); this gate is the
# single contract callers rely on.
wo = candidate.get("work_order_id")
if not wo or not _WO_ID_RE.match(str(wo)):
return False, "missing_required_field"
# email_type in enum. isinstance guard first: a non-str model value (JSON
# list/dict) is unhashable and `x in <set>` would raise, escaping the gate
# into async retries -- the opposite of fail-closed.
et = candidate.get("email_type")
if not isinstance(et, str) or et not in VALID_EMAIL_TYPES:
return False, "missing_required_field"
# status if non-null in the enum (same unhashable-type guard).
status = candidate.get("status")
if status is not None and (
not isinstance(status, str) or status not in VALID_STATUSES
):
return False, "invalid_status"
# site_code if non-null matches the code pattern
site = candidate.get("site_code")
if site is not None and not _SITE_CODE_RE.match(str(site)):
return False, "malformed_site_code"
# Date fields if non-null must be ISO-8601 strings (the prompt's declared
# format), converging with the template path's parsed-date guarantee and
# keeping the _UNPARSEABLE sentinel / arbitrary model prose out of the
# store.
for key in ("date_reported", "scheduled_start", "due_date", "comment_time"):
val = candidate.get(key)
if val is None:
continue
try:
datetime.fromisoformat(str(val))
except ValueError:
return False, "creation_time_unparseable"
# Free-text scalars must be None or str. A prompt-injected float reaches
# DynamoDB as Python float and crashes update_item; a dict/list lands as
# a Map/List attribute and pollutes downstream readers (PO parity).
for field in _AI_FREE_TEXT_STR_FIELDS:
val = candidate.get(field)
if val is not None and not isinstance(val, str):
return False, "invalid_field_type"
return True, "ok"
def try_deterministic_parse(email_data):
"""Entry point. Returns (parsed|None, parse_method, template_id, reason).
On a proven-conformant parse returns (dict, 'template', template_id, 'ok').
On any miss/invalid/exception returns (None, 'ai_fallback', template_id,
reason) -- a failure is NEVER a parsed result."""
template_id = "unknown"
try:
template_id, reason = classify_template(email_data)
if template_id == "unknown":
return None, "ai_fallback", template_id, reason
if template_id == "update_plaintext":
candidate = extract_update_plaintext(email_data)
else:
candidate = extract_assign_html(email_data)
ok, reason = validate(candidate, template_id, email_data)
if not ok:
return None, "ai_fallback", template_id, reason
return candidate, "template", template_id, "ok"
except Exception: # noqa: BLE001 -- fail closed on ANY extractor error
return None, "ai_fallback", template_id, "extractor_raised"