apm-wo-analysis/lambdas/classifier/classify.py
Adam Moussa eb724580e0 Add two-axis classifier Lambda and CDK wiring (Phase 2)
Implement the core classification engine and wire it into the pipeline stack.

classify.py: two-axis classifier — HTML-strip, comment-intent regex buckets
(escalations → status inquiry, most-specific first), Hold Reason / WO Status
structured state, comment-vs-state mismatch detector, and a Claude Haiku
fallback (Secrets Manager key) reserved for ambiguous free-text. Exports
ESCALATION_CATEGORIES / ACTION_NEEDED_CATEGORIES.

handler.py: S3-triggered handler — parse xlsx/csv, classify each non-blank
row, write a per-WO Parquet snapshot to analytics/dt=YYYY-MM-DD/ (registers
the Glue partition via awswrangler) and a summary.json for slack-post (Phase 4).

pipeline_stack.py: Glue database, ARM64 Python 3.12 classifier Lambda
(Docker-bundled deps), S3 raw/ notification (.xlsx/.csv), and least-privilege
IAM (read raw/, read-write analytics/, scoped Glue catalog, read Anthropic key).

Smoke-tested against the real export: 347 rows, "Other" at 5.2% (target ~9%),
18 mismatches flagged. 7/7 unit + smoke tests pass; cdk synth green.
2026-05-28 17:09:52 -04:00

551 lines
19 KiB
Python
Raw Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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.

"""Two-axis APM work-order classifier — the core domain logic.
Axis 1 — comment intent: regex over the HTML-stripped ``Last Comment``,
most-specific first.
Axis 2 — structured state: ``Hold Reason`` + ``WO Status``.
Resolution: comment intent wins when confident, else structured state, else
``Other``. A Claude Haiku fallback (Secrets Manager
``apm-wo-analysis/anthropic-api-key``) is reserved strictly for ambiguous
free-text with no structured signal. A mismatch detector flags when comment
intent contradicts structured state — surface, never suppress.
The authoritative spec is CLAUDE.md ("The classification model"). This module
is owned by the classifier-engineer agent; the rule ladder and Haiku fallback
are implemented in Phase 2 (docs/BUILD.md) and smoke-tested against a real
export before merge.
``classify()`` is fully deterministic: no network, no AWS, no Secrets access.
The Haiku fallback lives in ``classify_with_haiku()`` and is never reached from
inside ``classify()``.
"""
from __future__ import annotations
import html
import os
import re
# Axis 2 — Hold Reason → category.
HOLD_TO_CATEGORY = {
"SCHEDULING": "Awaiting Scheduling",
"REPORT": "Report / Docs Needed",
"VENDOR": "Awaiting Vendor / Parts",
"PARTS": "Awaiting Vendor / Parts",
"ORDER": "Awaiting Vendor / Parts",
"VERIFY": "Verification Needed",
"RESOURCE": "Resource Hold",
"NOEQUIP": "No Equipment",
}
# Axis 2 — WO Status signals.
WO_STATUS_CANCELLED = "RCAN" # → Cancelled
WO_STATUS_HOLD = "H" # corroborates On Hold
WO_STATUS_IN_FLIGHT = frozenset({"IP", "R", "RR"})
ESCALATION_CATEGORIES = frozenset(
{
"1st Escalation",
"2nd Escalation",
"3rd Escalation",
"SIM Ticket",
"Other Escalation",
}
)
ACTION_NEEDED_CATEGORIES = ESCALATION_CATEGORIES | frozenset(
{
"Awaiting Scheduling",
"Report / Docs Needed",
"Awaiting Report / Invoice",
"Awaiting Vendor / Parts",
"Status Inquiry",
"Vendor No-Show",
}
)
# ---------------------------------------------------------------------------
# Axis 1 — HTML stripping
# ---------------------------------------------------------------------------
_TAG_RE = re.compile(r"<[^>]+>")
_URL_RE = re.compile(r"https?://\S+", re.IGNORECASE)
_WS_RE = re.compile(r"\s+")
# Smart quotes / dashes → ASCII so the rule ladder can use plain patterns.
_SMART_MAP = {
"‘": "'", # left single quote
"’": "'", # right single quote
"‚": "'", # single low-9 quote
"‛": "'", # single high-reversed-9 quote
"“": '"', # left double quote
"”": '"', # right double quote
"„": '"', # double low-9 quote
"«": '"', # left guillemet
"»": '"', # right guillemet
"–": "-", # en dash
"—": "-", # em dash
"…": "...", # ellipsis
" ": " ", # non-breaking space
}
_SMART_RE = re.compile("|".join(map(re.escape, _SMART_MAP)))
def strip_html(comment: str | None) -> str:
"""Unwrap the HTML-wrapped Last Comment into clean plain text.
Drops the ``<html>...</html>`` wrapper and all nested tags, decodes HTML
entities, normalises smart quotes/dashes to ASCII, reduces embedded URLs to
a readable ``link`` token, and collapses whitespace. ``None``/blank → "".
"""
if comment is None:
return ""
text = str(comment)
# Decode entities first so e.g. &lt;div&gt; becomes a real tag we can strip,
# then strip tags, then decode again for entities that were *inside* text.
text = html.unescape(text)
text = _TAG_RE.sub(" ", text)
text = html.unescape(text)
# Reduce URLs to a readable token (the raw link is noise for matching).
text = _URL_RE.sub(" link ", text)
text = _SMART_RE.sub(lambda m: _SMART_MAP[m.group(0)], text)
text = _WS_RE.sub(" ", text).strip()
return text
# ---------------------------------------------------------------------------
# Axis 1 — comment intent rule ladder (most-specific first)
# ---------------------------------------------------------------------------
def comment_intent(comment: str) -> str | None:
"""Ordered, most-specific-first rule ladder over the *stripped* text.
Returns one of the Axis-1 buckets in CLAUDE.md, or ``None`` when no rule
fires. The caller is responsible for stripping HTML first; this operates on
plain text. Ordering matters: the most specific / highest-severity intents
are tested before generic ones so e.g. a 3rd-attempt escalation never falls
through to a generic "report needed" rule.
"""
if not comment:
return None
t = comment
low = t.lower()
# -- Escalations: APM "Nth attempt process" cadence + explicit escalations.
# Most-specific first: 3rd → 2nd → 1st. "Email/phone escalation sent" is the
# escalation tell; the "attempt" number sets the level. A bare 1st-attempt
# schedule-confirmation request is routine outreach, not an escalation, so
# require an escalation/attempt signal beyond the plain "1st attempt".
if re.search(r"\b3rd\b.*\battempt\b", low) or "3rd escalation" in low:
return "3rd Escalation"
if re.search(r"\b2nd\b.*\battempt\b", low) or "2nd escalation" in low:
return "2nd Escalation"
if "1st escalation" in low:
return "1st Escalation"
# -- SIM Ticket: Amazon SIM tooling references.
if (
re.search(r"\bsim\b.*\b(ticket|tt|v\d{6,})\b", low)
or "t.corp.amazon.com" in low
or re.search(r"\bsim\s*tt\b", low)
):
return "SIM Ticket"
# -- Vendor No-Show: vendor failed to appear for a scheduled visit.
if (
"no show" in low
or "no-show" in low
or "noshow" in low
or "did not show up" in low
or "didn't show up" in low
or "did not show" in low
or "failed to show" in low
):
return "Vendor No-Show"
# -- Cancelled: explicit cancellation language in the comment.
if re.search(r"\bwo cancelled\b", low) or re.search(
r"\bcancel(l)?ed\b.*\b(error|deprecat)", low
):
return "Cancelled"
# -- Weekly WO Scheduled: the recurring-weekly cadence template.
if re.search(r"\bweekly\s+wo\s+scheduled\b", low) or (
"weekly" in low and "service reports are required weekly" in low
):
return "Weekly WO Scheduled"
# -- Avetta Project Created: an Avetta project/work-request was opened.
if re.search(r"\bavetta\s+project\s+created\b", low) or re.search(
r"\bwork request\b.*\bcreated in avetta\b", low
):
return "Avetta Project Created"
# -- Schedule Confirmed: vendor/site has a confirmed (re)schedule.
if (
"schedule confirmed with vendor" in low
or "scheduled confirmed with vendor" in low
or re.search(r"\bschedule is confirmed\b", low)
or re.search(r"\b(re)?scheduling confirmed\b", low)
or re.search(r"\bconfirmed scheduled date\b", low)
or re.search(r"\bconfirmed schedule(d)? (date|for)\b", low)
or re.search(r"\bvendor confirmed\b", low)
or re.search(r"^confirmed\b", low)
):
return "Schedule Confirmed"
# -- Completed / Pending Close: work done but WO not yet closed.
if (
re.search(r"\bperformed task\b", low)
or re.search(r"\bcompleted by\b", low)
or re.search(r"\bcan'?t close\b", low) # cant/can't close
or re.search(r"\bcannot close\b", low)
or re.search(r"\bunable to close\b", low)
or "will not allow me to close" in low
or "will not let me close" in low
or re.search(r"\bpending close\b", low)
):
return "Completed / Pending Close"
# -- Rescheduled: an explicit reschedule/new-date event (no confirmation).
if (
re.search(r"\brescheduled\b", low)
or re.search(r"\bre-?scheduling\b", low)
or re.search(r"\bnew eta\b", low)
or re.search(r"\btentative(ly)? schedule", low)
):
return "Rescheduled"
# -- Report / Docs Needed: completion/service report or docs are requested.
if (
"completion report" in low
or "service report" in low
or "completion validation" in low
or re.search(r"\bservice reports? (are )?(required|needed)\b", low)
or re.search(r"\bpm service report\b", low)
or re.search(r"\b(upload|provide).{0,30}\b(report|document|docs|jha)\b", low)
or re.search(r"\bclosing documents\b", low)
or re.search(r"\bcompletion document", low)
):
return "Report / Docs Needed"
# -- Awaiting Report / Invoice: waiting on a report or invoice to land.
if (
re.search(r"\bpending reports?\b", low)
or re.search(r"\bawaiting (the )?(report|invoice)\b", low)
or re.search(r"\binvoice from\b", low)
or re.search(r"\bawaiting invoice\b", low)
):
return "Awaiting Report / Invoice"
# -- Awaiting Scheduling: asking the team/vendor to schedule or confirm.
if (
re.search(r"\bplease s?c?hedule\b", low) # "please schedule"/"shedule" typo
or re.search(r"\bplease provide schedul", low)
or re.search(r"\bplease advise on schedul", low)
or re.search(r"\bprovide update on schedul", low)
or re.search(r"\bneeds? (a )?confirmed schedul", low)
or re.search(r"\bget (this )?scheduled\b", low)
or re.search(r"\bschedule confirmation\b", low)
or re.search(r"\bupdates? on scheduling\b", low)
or re.search(r"\bget this scheduled\b", low)
or re.search(r"\bwaiting for a schedule", low)
or re.search(r"\bawaiting (a )?schedul", low)
or re.search(r"\b1\s*st attempt process for schedule confirmation\b", low)
or re.search(r"\brequesting service date\b", low)
or re.search(r"\bneed(s|ed)? to (be )?(re)?schedul", low)
or re.search(r"\bcan we (please )?get this scheduled\b", low)
):
return "Awaiting Scheduling"
# -- Status Inquiry: a question about status/ETA/updates.
if (
re.search(r"\bany update", low)
or re.search(r"\bupdates?\?", low)
or re.search(r"\beta\b", low) # any mention of an ETA is a status inquiry
or re.search(r"\bdo you have an eta\b", low)
or (low.endswith("?") and ("update" in low or "eta" in low))
or re.search(r"\bplease confirm whether\b", low)
or re.search(r"\bconfirm completion or scheduling status\b", low)
):
return "Status Inquiry"
# -- On Hold: explicit on-hold language in the comment.
if re.search(r"\bon[- ]?hold\b", low) or re.search(r"\bwork order on hold\b", low):
return "On Hold"
# -- Acknowledgement / No-op: terse acks ("Copy", "Copy 5/26.").
if re.fullmatch(r"copy[.\s\d/]*", low) or re.fullmatch(
r"(noted|ack(nowledged)?)[.\s]*", low
):
return "Acknowledgement / No-op"
return None
# ---------------------------------------------------------------------------
# Resolution policy + mismatch detection
# ---------------------------------------------------------------------------
# Comment intents that assert the work is essentially done / scheduled.
_DONE_INTENTS = frozenset(
{
"Completed / Pending Close",
"Schedule Confirmed",
}
)
def _hold_category(hold_reason: str | None) -> str | None:
if not hold_reason:
return None
return HOLD_TO_CATEGORY.get(hold_reason.strip().upper())
def _mismatch(
intent: str | None,
wo_status: str | None,
hold_reason: str | None,
) -> str | None:
"""Flag when comment intent contradicts the structured state.
A mismatch is a feature, not noise: it surfaces WOs whose free-text claims
progress the structured state has not caught up with. The returned string is
human-readable and goes straight onto the snapshot row.
"""
if intent is None:
return None
status = (wo_status or "").strip().upper()
hold = (hold_reason or "").strip().upper()
# Comment claims completion/confirmation while still on a blocking hold.
if intent in _DONE_INTENTS and hold in {
"REPORT",
"SCHEDULING",
"VENDOR",
"PARTS",
"ORDER",
}:
verb = (
"completed"
if intent == "Completed / Pending Close"
else "schedule confirmed"
)
return f"comment says {verb} but WO is on a {hold} hold"
# Comment claims completion while the WO Status is still in-flight.
if intent == "Completed / Pending Close" and status in WO_STATUS_IN_FLIGHT:
return f"comment says completed but WO Status is {status}"
# Comment claims a vendor no-show while the WO carries a schedule-confirmed
# type hold — the two stories disagree.
if intent == "Vendor No-Show" and hold == "SCHEDULING":
return "comment reports vendor no-show but WO is on a SCHEDULING hold"
return None
def classify(
wo_status: str | None,
hold_reason: str | None,
last_comment: str | None,
) -> tuple[str, str | None]:
"""Resolve to ``(final_category, mismatch_reason | None)``.
Fully deterministic — no network, no AWS. Resolution policy:
1. Comment intent wins when confident (the ladder fired).
2. Else structured state: ``Hold Reason`` map, then ``WO Status``
(``RCAN`` → Cancelled, ``H`` corroborates On Hold).
3. Else ``Other``.
Regardless of which axis decides the category, a mismatch reason is computed
whenever the comment intent contradicts the structured state and surfaced on
the result (never suppressed).
"""
status = (wo_status or "").strip().upper()
text = strip_html(last_comment)
intent = comment_intent(text)
mismatch = _mismatch(intent, status, hold_reason)
# 1. Comment intent wins when confident.
if intent is not None:
return intent, mismatch
# 2. Structured state.
hold_cat = _hold_category(hold_reason)
if hold_cat is not None:
return hold_cat, mismatch
if status == WO_STATUS_CANCELLED:
return "Cancelled", mismatch
if status == WO_STATUS_HOLD:
return "On Hold", mismatch
# 3. Nothing matched.
return "Other", mismatch
# ---------------------------------------------------------------------------
# Haiku fallback — isolated, lazy, network-only, NEVER reached from classify()
# ---------------------------------------------------------------------------
_HAIKU_MODEL = "claude-haiku-4-5-20251001"
_SECRET_NAME = "apm-wo-analysis/anthropic-api-key"
_ANTHROPIC_URL = "https://api.anthropic.com/v1/messages"
_ANTHROPIC_VERSION = "2023-06-01"
# Buckets Haiku is allowed to choose from (Axis-1 taxonomy).
_HAIKU_BUCKETS = (
"1st Escalation",
"2nd Escalation",
"3rd Escalation",
"SIM Ticket",
"Vendor No-Show",
"Weekly WO Scheduled",
"Schedule Confirmed",
"Awaiting Scheduling",
"Report / Docs Needed",
"Awaiting Report / Invoice",
"Completed / Pending Close",
"Acknowledgement / No-op",
"Cancelled",
"On Hold",
"Rescheduled",
"Avetta Project Created",
"Other Escalation",
"Status Inquiry",
"Other",
)
def _haiku_enabled() -> bool:
"""The fallback is gated behind ``APM_HAIKU_FALLBACK`` (default enabled)."""
return os.environ.get("APM_HAIKU_FALLBACK", "1").strip().lower() not in {
"0",
"false",
"no",
"off",
"",
}
def _is_trivial(text: str) -> bool:
"""Too short / contentless to be worth an LLM call."""
return len(text) < 4
def _fetch_api_key() -> str:
"""Read the Anthropic key from Secrets Manager (lazy, call-time only)."""
import json
import boto3 # imported lazily so classify()/tests never need boto3
client = boto3.client("secretsmanager")
resp = client.get_secret_value(SecretId=_SECRET_NAME)
raw = resp.get("SecretString") or ""
raw = raw.strip()
# Support both a bare string secret and a JSON {"anthropic-api-key": "..."}.
if raw.startswith("{"):
data = json.loads(raw)
for key in ("anthropic-api-key", "ANTHROPIC_API_KEY", "api_key", "key"):
if key in data:
return data[key]
# Single-value JSON object: take the only value.
if len(data) == 1:
return next(iter(data.values()))
raise KeyError("anthropic api key not found in secret JSON")
return raw
def _call_haiku(text: str, api_key: str) -> str | None:
"""POST to the Anthropic messages API via stdlib urllib; return a bucket."""
import json
import urllib.error
import urllib.request
bucket_list = "\n".join(f"- {b}" for b in _HAIKU_BUCKETS)
prompt = (
"You classify a maintenance work-order's last comment into exactly one "
"category. The deterministic rules failed and there is no structured "
"hold signal, so judge the free text alone. Reply with ONLY the exact "
"category name from this list, nothing else:\n"
f"{bucket_list}\n\n"
f"Comment:\n{text}"
)
body = json.dumps(
{
"model": _HAIKU_MODEL,
"max_tokens": 20,
"messages": [{"role": "user", "content": prompt}],
}
).encode("utf-8")
req = urllib.request.Request(
_ANTHROPIC_URL,
data=body,
headers={
"content-type": "application/json",
"x-api-key": api_key,
"anthropic-version": _ANTHROPIC_VERSION,
},
method="POST",
)
try:
with urllib.request.urlopen(req, timeout=15) as resp:
payload = json.loads(resp.read().decode("utf-8"))
except (urllib.error.URLError, TimeoutError, ValueError):
return None
parts = payload.get("content") or []
raw = "".join(p.get("text", "") for p in parts if isinstance(p, dict)).strip()
if not raw:
return None
# Normalise: pick the listed bucket that matches the model's reply.
for bucket in _HAIKU_BUCKETS:
if bucket.lower() == raw.lower():
return bucket
for bucket in _HAIKU_BUCKETS:
if bucket.lower() in raw.lower():
return bucket
return None
def classify_with_haiku(
wo_status: str | None,
hold_reason: str | None,
last_comment: str | None,
) -> tuple[str, str | None]:
"""Deterministic ``classify()`` first; Haiku only for true residual ``Other``.
The LLM is consulted strictly when ALL hold:
* deterministic result is ``Other``;
* ``Hold Reason`` is blank (no structured signal to fall back on);
* the stripped comment is non-trivial;
* ``APM_HAIKU_FALLBACK`` is enabled.
Everything network-related (boto3, the Anthropic call) is imported and
invoked lazily here, so ``classify()`` and the test-suite never touch it.
The mismatch reason from the deterministic pass is preserved.
"""
category, mismatch = classify(wo_status, hold_reason, last_comment)
if category != "Other":
return category, mismatch
if (hold_reason or "").strip():
return category, mismatch
if not _haiku_enabled():
return category, mismatch
text = strip_html(last_comment)
if _is_trivial(text):
return category, mismatch
try:
api_key = _fetch_api_key()
guess = _call_haiku(text, api_key)
except Exception:
return category, mismatch
if guess and guess != "Other":
return guess, mismatch
return category, mismatch