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.
This commit is contained in:
Adam Moussa 2026-05-28 17:09:52 -04:00
parent 7befc8f0d7
commit eb724580e0
4 changed files with 1045 additions and 32 deletions

View file

@ -1,27 +1,50 @@
"""Pipeline stack: S3, classifier + slack-post Lambdas, Glue, Athena, IAM.
Scaffold — the exports bucket (Phase 1 of docs/BUILD.md) is included so the
stack synthesizes to something real. The classifier Lambda (Phase 2), Glue
database + Athena workgroup with partition projection (Phase 3), and the
slack-post Lambda + IAM (Phase 4) are added in their respective phases.
Phase 1 — exports bucket + drop-folder uploader IAM user.
Phase 2 — Glue database + S3-triggered classifier Lambda (this file).
Phase 3 — Athena workgroup + partition-projection table (TODO).
Phase 4 — slack-post Lambda + scoped IAM (TODO).
Lambda defaults when added: Python 3.12, ARM64, explicit LogGroup with 60-day
retention. No DynamoDB — this is an S3 + Athena analytics workload (see CLAUDE.md).
Lambdas: Python 3.12, ARM64, explicit LogGroup with 60-day retention.
No DynamoDB — this is an S3 + Athena analytics workload (see CLAUDE.md).
"""
import os
from aws_cdk import (
BundlingOptions,
Duration,
RemovalPolicy,
Stack,
)
from aws_cdk import (
aws_glue as glue,
)
from aws_cdk import (
aws_iam as iam,
)
from aws_cdk import (
aws_lambda as lambda_,
)
from aws_cdk import (
aws_logs as logs,
)
from aws_cdk import (
aws_s3 as s3,
)
from aws_cdk import (
aws_s3_notifications as s3n,
)
from aws_cdk import (
aws_secretsmanager as secretsmanager,
)
from constructs import Construct
GLUE_DATABASE = "apm_wo_analysis"
GLUE_TABLE = "apm_wo_snapshots"
ANTHROPIC_SECRET = "apm-wo-analysis/anthropic-api-key"
LAMBDAS_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "lambdas")
class PipelineStack(Stack):
def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None:
@ -61,7 +84,87 @@ class PipelineStack(Stack):
)
)
# Phase 2 — classifier Lambda, S3-triggered on the raw/ prefix. TODO
# Phase 3 — Glue database `apm_wo_analysis` + Athena workgroup
# (partition projection on dt; no crawler). TODO
# Phase 4 — slack-post Lambda + scoped IAM. TODO
# Phase 2/3 — Glue database for the analytics dataset. The classifier
# registers the apm_wo_snapshots table/partitions into it via awswrangler;
# Phase 3 swaps in the partition-projection table definition + Athena.
glue.CfnDatabase(
self,
"AnalyticsDb",
catalog_id=self.account,
database_input=glue.CfnDatabase.DatabaseInputProperty(name=GLUE_DATABASE),
)
# Phase 2 — classifier Lambda, S3-triggered on the raw/ prefix.
# Deps (awswrangler/pandas/pyarrow/openpyxl) are Docker-bundled for ARM64
# from lambdas/classifier/requirements.txt; boto3 ships in the runtime.
classifier_logs = logs.LogGroup(
self,
"ClassifierLogs",
log_group_name="/aws/lambda/apm-wo-analysis-classifier",
retention=logs.RetentionDays.TWO_MONTHS,
removal_policy=RemovalPolicy.DESTROY,
)
self.classifier_fn = lambda_.Function(
self,
"Classifier",
function_name="apm-wo-analysis-classifier",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
memory_size=512,
timeout=Duration.seconds(120),
log_group=classifier_logs,
environment={"APM_HAIKU_FALLBACK": "on"},
code=lambda_.Code.from_asset(
os.path.join(LAMBDAS_DIR, "classifier"),
bundling=BundlingOptions(
image=lambda_.Runtime.PYTHON_3_12.bundling_image,
platform="linux/arm64",
command=[
"bash",
"-c",
"pip install -r requirements.txt -t /asset-output "
"&& cp -au . /asset-output",
],
),
),
)
# S3 trigger: any .xlsx/.csv landing under raw/ invokes the classifier.
for suffix in (".xlsx", ".csv"):
self.exports_bucket.add_event_notification(
s3.EventType.OBJECT_CREATED,
s3n.LambdaDestination(self.classifier_fn),
s3.NotificationKeyFilter(prefix="raw/", suffix=suffix),
)
# IAM — least privilege: read raw/, read+write analytics/, catalog the
# snapshots table, and read the Anthropic key for the Haiku fallback.
self.exports_bucket.grant_read(self.classifier_fn, "raw/*")
self.exports_bucket.grant_read_write(self.classifier_fn, "analytics/*")
self.classifier_fn.add_to_role_policy(
iam.PolicyStatement(
sid="GlueCatalogSnapshots",
actions=[
"glue:GetDatabase",
"glue:GetTable",
"glue:CreateTable",
"glue:UpdateTable",
"glue:GetPartition",
"glue:GetPartitions",
"glue:CreatePartition",
"glue:BatchCreatePartition",
"glue:UpdatePartition",
],
resources=[
f"arn:aws:glue:{self.region}:{self.account}:catalog",
f"arn:aws:glue:{self.region}:{self.account}:database/{GLUE_DATABASE}",
f"arn:aws:glue:{self.region}:{self.account}:table/{GLUE_DATABASE}/*",
],
)
)
secretsmanager.Secret.from_secret_name_v2(
self, "AnthropicKey", ANTHROPIC_SECRET
).grant_read(self.classifier_fn)
# Phase 4 — slack-post Lambda + scoped IAM; classifier async-invokes it. TODO

View file

@ -14,10 +14,18 @@ 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",
@ -56,15 +64,286 @@ ACTION_NEEDED_CATEGORIES = ESCALATION_CATEGORIES | frozenset(
}
)
# ---------------------------------------------------------------------------
# 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 and decode entities. (Phase 2)"""
raise NotImplementedError
"""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; returns a bucket or None. (Phase 2)"""
raise NotImplementedError
"""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(
@ -72,5 +351,201 @@ def classify(
hold_reason: str | None,
last_comment: str | None,
) -> tuple[str, str | None]:
"""Resolve to ``(final_category, mismatch_reason | None)``. (Phase 2)"""
raise NotImplementedError
"""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

View file

@ -1,19 +1,232 @@
"""S3-triggered classifier Lambda: parse export → two-axis classify → Parquet.
"""S3-triggered classifier Lambda: parse export -> two-axis classify -> Parquet.
Triggered on ``s3:ObjectCreated`` under the ``raw/`` prefix regardless of how
the file arrives (direct upload or local drop-folder). For each non-blank-comment
row it strips HTML, runs ``classify()``, derives ``is_escalation`` / ``is_action``
/ ``mismatch``, writes a per-WO snapshot to ``analytics/dt=YYYY-MM-DD/`` as
Parquet (registering the Glue partition), emits a small ``summary.json`` for the
slack-post Lambda, then invokes it.
Triggered on ``s3:ObjectCreated`` under the ``raw/`` prefix regardless of how the
file arrives (direct upload or local drop-folder). For each non-blank-comment row
it strips HTML, runs the two-axis classifier, derives ``is_escalation`` /
``is_action`` / ``mismatch``, writes a per-WO snapshot to ``analytics/dt=YYYY-MM-DD/``
as Parquet (registering the Glue partition), and emits a small ``summary.json``
for the slack-post Lambda to read cheaply.
Runtime when wired up: Python 3.12, ARM64, 512 MB, 120 s, awswrangler layer.
Implemented in Phase 2 (docs/BUILD.md).
The slack-post invocation is wired in Phase 4 (the function does not exist yet).
Runtime: Python 3.12, ARM64, 512 MB, 120 s. pandas/pyarrow/awswrangler are
bundled from requirements.txt; boto3 is provided by the Lambda runtime.
"""
from __future__ import annotations
import csv
import json
import os
import tempfile
from collections import Counter
from datetime import datetime, timezone
from urllib.parse import unquote_plus
import awswrangler as wr
import boto3
import openpyxl
import pandas as pd
import classify as clf
GLUE_DATABASE = "apm_wo_analysis"
GLUE_TABLE = "apm_wo_snapshots"
ANALYTICS_PREFIX = "analytics"
# Map snapshot field -> substring matched (case-insensitively) against the export
# header, tolerating minor header drift in the 13-column APM export.
COLUMN_MATCHERS = {
"wo_number": "wo number",
"wo_description": "wo description",
"equipment_code": "equipment",
"site": "organization",
"due_date": "due date",
"department": "department",
"wo_status": "wo status",
"hold_reason": "hold reason",
"last_comment": "last comment",
"last_comment_by": "last comment by",
"last_comment_date": "last comment date",
"contractor": "contractor",
"contractor_description": "contractor description",
}
_s3 = boto3.client("s3")
def _haiku_enabled() -> bool:
return os.environ.get("APM_HAIKU_FALLBACK", "on").strip().lower() in (
"1",
"on",
"true",
"yes",
)
def _resolve_classifier():
"""Use the Haiku-fallback classifier when enabled, else the deterministic one."""
if _haiku_enabled() and hasattr(clf, "classify_with_haiku"):
return clf.classify_with_haiku
return clf.classify
def _norm(value) -> str:
if value is None:
return ""
if isinstance(value, datetime):
return value.isoformat()
return str(value).strip()
def _resolve_columns(header: list[str]) -> dict[str, int]:
"""Map each snapshot field to its column index in the export header."""
lowered = [(_norm(h).lower(), i) for i, h in enumerate(header)]
resolved: dict[str, int] = {}
for field, needle in COLUMN_MATCHERS.items():
# "last comment" is a prefix of "last comment by"/"date"; prefer exact-ish.
exact = [i for h, i in lowered if h == needle]
contains = [i for h, i in lowered if needle in h]
match = exact or contains
if match:
resolved[field] = match[0]
return resolved
def _read_rows(path: str, key: str) -> tuple[list[str], list[list]]:
"""Return (header, data_rows) from an xlsx or csv export."""
if key.lower().endswith(".csv"):
with open(path, newline="", encoding="utf-8-sig") as fh:
reader = list(csv.reader(fh))
return reader[0], reader[1:]
wb = openpyxl.load_workbook(path, read_only=True, data_only=True)
ws = wb.active
rows = list(ws.iter_rows(values_only=True))
wb.close()
return list(rows[0]), [list(r) for r in rows[1:]]
def _build_snapshot(header: list[str], data: list[list]) -> tuple[pd.DataFrame, int]:
"""Classify each non-blank-comment row into a snapshot DataFrame."""
cols = _resolve_columns(header)
if "last_comment" not in cols or "wo_status" not in cols:
raise ValueError(f"Export missing required columns; resolved={list(cols)}")
classifier = _resolve_classifier()
records: list[dict] = []
blank = 0
for row in data:
def cell(field: str):
idx = cols.get(field)
return row[idx] if idx is not None and idx < len(row) else None
raw_comment = cell("last_comment")
comment = clf.strip_html(raw_comment)
if not comment:
blank += 1 # blank-comment rows excluded from the classified total
continue
status = _norm(cell("wo_status"))
hold = _norm(cell("hold_reason"))
category, mismatch = classifier(status, hold, comment)
records.append(
{
"wo_number": _norm(cell("wo_number")),
"wo_description": _norm(cell("wo_description")),
"equipment_code": _norm(cell("equipment_code")),
"site": _norm(cell("site")),
"due_date": _norm(cell("due_date")),
"department": _norm(cell("department")),
"wo_status": status,
"hold_reason": hold,
"last_comment": comment,
"last_comment_by": _norm(cell("last_comment_by")),
"last_comment_date": _norm(cell("last_comment_date")),
"contractor": _norm(cell("contractor")),
"contractor_description": _norm(cell("contractor_description")),
"category": category,
"is_escalation": category in clf.ESCALATION_CATEGORIES,
"is_action": category in clf.ACTION_NEEDED_CATEGORIES,
"mismatch": mismatch or "",
}
)
return pd.DataFrame.from_records(records), blank
def _build_summary(df: pd.DataFrame, dt: str, key: str, blank: int) -> dict:
category_counts = Counter(df["category"])
site_counts = Counter(s for s in df["site"] if s)
mismatches = [
{"wo_number": r.wo_number, "category": r.category, "mismatch": r.mismatch}
for r in df.itertuples()
if r.mismatch
]
return {
"dt": dt,
"source_key": key,
"classified_total": int(len(df)),
"blank_comment_rows": int(blank),
"category_counts": dict(category_counts),
"escalation_total": int(df["is_escalation"].sum()),
"third_escalation_count": int(category_counts.get("3rd Escalation", 0)),
"action_needed": int(df["is_action"].sum()),
"routine": int((~df["is_action"]).sum()),
"top_sites": [{"site": s, "count": n} for s, n in site_counts.most_common(10)],
"mismatches": mismatches,
"generated_at": datetime.now(timezone.utc).isoformat(),
}
def handler(event, context):
"""Lambda entry point. (Phase 2)"""
raise NotImplementedError
"""Classify each export dropped under raw/ into a daily Parquet snapshot."""
record = event["Records"][0]
bucket = record["s3"]["bucket"]["name"]
key = unquote_plus(record["s3"]["object"]["key"])
if not key.startswith("raw/") or not key.lower().endswith((".xlsx", ".csv")):
print(f"Skipping non-export object: s3://{bucket}/{key}")
return {"skipped": key}
dt = datetime.now(timezone.utc).strftime("%Y-%m-%d")
print(f"Classifying s3://{bucket}/{key} into dt={dt}")
with tempfile.NamedTemporaryFile(suffix=os.path.splitext(key)[1]) as tmp:
_s3.download_fileobj(bucket, key, tmp)
tmp.flush()
header, data = _read_rows(tmp.name, key)
df, blank = _build_snapshot(header, data)
if df.empty:
print("No classifiable rows (all comments blank); nothing written.")
return {"classified": 0, "blank": blank}
df["dt"] = dt
wr.s3.to_parquet(
df=df,
path=f"s3://{bucket}/{ANALYTICS_PREFIX}/",
dataset=True,
partition_cols=["dt"],
mode="overwrite_partitions",
database=GLUE_DATABASE,
table=GLUE_TABLE,
)
summary = _build_summary(df, dt, key, blank)
_s3.put_object(
Bucket=bucket,
Key=f"{ANALYTICS_PREFIX}/dt={dt}/summary.json",
Body=json.dumps(summary, indent=2).encode("utf-8"),
ContentType="application/json",
)
# Phase 4: async-invoke the slack-post Lambda here once it exists.
print(
f"Wrote {len(df)} rows, {summary['escalation_total']} escalations "
f"({summary['third_escalation_count']} 3rd), {len(summary['mismatches'])} mismatches."
)
return {
"classified": int(len(df)),
"dt": dt,
"summary_key": f"dt={dt}/summary.json",
}

View file

@ -1,24 +1,246 @@
"""Smoke test for the two-axis classifier.
The full smoke test — run ``classify()`` over the sample export
(~/Downloads/_documents/Sheet1-1.xlsx) and assert "Other" <= ~10% with known
fixtures landing in the right buckets — is added in Phase 2 (docs/BUILD.md) and
owned by the classifier-engineer agent. This placeholder keeps the test layout
in place and guards the classification constants.
Runs the DETERMINISTIC ``classify()`` (no Haiku, no network, no AWS) over the
canonical sample export and asserts the two-axis model keeps "Other" in the
single digits, plus a handful of known fixtures land in the right bucket. The
Haiku fallback is never exercised here.
Canonical fixture: ``~/Downloads/_documents/Sheet1-1.xlsx`` (347 rows, 13 cols).
If the file is absent the export-driven tests skip. Run with the repo venv:
./.venv/bin/python -m pytest tests/test_classify.py -s -q
"""
import collections
import sys
from pathlib import Path
import pytest
sys.path.insert(
0, str(Path(__file__).resolve().parent.parent / "lambdas" / "classifier")
)
import classify # noqa: E402
# Column indices in the 13-column export.
COL_WO_STATUS = 6
COL_HOLD_REASON = 7
COL_LAST_COMMENT = 8
FIXTURE = Path.home() / "Downloads" / "_documents" / "Sheet1-1.xlsx"
# Max acceptable deterministic "Other" share before any AI. CLAUDE.md: the
# two-axis model lands ~9% before Haiku; we hold the line at single digits.
MAX_OTHER_PCT = 10.0
def test_classification_constants_well_formed():
assert "3rd Escalation" in classify.ESCALATION_CATEGORIES
assert classify.HOLD_TO_CATEGORY["SCHEDULING"] == "Awaiting Scheduling"
# Every escalation category is also action-needed.
assert classify.ESCALATION_CATEGORIES <= classify.ACTION_NEEDED_CATEGORIES
def test_strip_html_unwraps_and_normalises():
assert classify.strip_html(None) == ""
assert classify.strip_html(" ") == ""
assert classify.strip_html("<html>WO schedule confirmed with vendor</html>") == (
"WO schedule confirmed with vendor"
)
# Nested tags, entities, smart quotes, and embedded URLs.
raw = (
"<html><div>1st attempt process for schedule confirmation. "
"Vendor, please confirm ‘Schedule Start Date’ &amp; proceed. "
"https://app.avetta.com/avt-cli/x</div></html>"
)
cleaned = classify.strip_html(raw)
assert "<" not in cleaned and ">" not in cleaned
assert "‘" not in cleaned and "’" not in cleaned # smart quotes gone
assert "'Schedule Start Date'" in cleaned
assert "&amp;" not in cleaned and "&" in cleaned # entity decoded
assert "https://" not in cleaned # URL reduced to a token
def test_known_intents():
# Schedule confirmed (the dominant happy-path comment).
cat, _ = classify.classify(
"IP", "", "<html>WO schedule confirmed with vendor</html>"
)
assert cat == "Schedule Confirmed"
# Cancelled via WO Status, even with a generic comment.
cat, _ = classify.classify(
"RCAN", "", "<html>WO Cancelled, created in error.</html>"
)
assert cat == "Cancelled"
# 3rd-attempt escalation cadence.
cat, _ = classify.classify(
"H", "REPORT", "<html>3rd attempt process for schedule confirmation.</html>"
)
assert cat == "3rd Escalation"
# Vendor no-show.
cat, _ = classify.classify(
"R", "", "<html>Site tech reported vendor was a no show for Friday.</html>"
)
assert cat == "Vendor No-Show"
# Weekly cadence template.
cat, _ = classify.classify(
"IP",
"",
"<html>Weekly WO scheduled. Service reports required EOD Friday.</html>",
)
assert cat == "Weekly WO Scheduled"
# Structured-only fallback: no comment intent fires, REPORT hold decides.
cat, _ = classify.classify("IP", "REPORT", "<html></html>")
assert cat == "Report / Docs Needed"
def test_mismatch_detection():
# Comment claims completion while on a REPORT hold → mismatch surfaced.
cat, mm = classify.classify(
"IP", "REPORT", "<html>Vendor arrived and performed task.</html>"
)
assert cat == "Completed / Pending Close"
assert mm is not None and "REPORT" in mm
# Schedule confirmed while on a SCHEDULING hold → mismatch surfaced.
cat, mm = classify.classify(
"R", "SCHEDULING", "<html>WO schedule confirmed with vendor.</html>"
)
assert cat == "Schedule Confirmed"
assert mm is not None
# Clean case: no contradiction → no mismatch.
_, mm = classify.classify(
"IP", "", "<html>WO schedule confirmed with vendor</html>"
)
assert mm is None
def test_classify_is_offline():
"""classify() must not import boto3 or reach the network."""
import sys as _sys
had_boto3 = "boto3" in _sys.modules
classify.classify("IP", "REPORT", "<html>1st attempt process for report.</html>")
# If boto3 wasn't already loaded, classify() must not have pulled it in.
if not had_boto3:
assert "boto3" not in _sys.modules
@pytest.mark.skipif(not FIXTURE.exists(), reason=f"sample export not found: {FIXTURE}")
def test_other_share_against_real_export(capsys):
import openpyxl
wb = openpyxl.load_workbook(FIXTURE, read_only=True, data_only=True)
rows = list(wb.active.iter_rows(values_only=True))[1:] # drop header
dist = collections.Counter()
mismatches = 0
classified = 0
blank = 0
other_samples = []
for row in rows:
comment = row[COL_LAST_COMMENT]
# Blank-comment rows are excluded from the classified total by design.
if comment is None or str(comment).strip() == "":
blank += 1
continue
classified += 1
category, mismatch = classify.classify(
row[COL_WO_STATUS], row[COL_HOLD_REASON], comment
)
dist[category] += 1
if mismatch:
mismatches += 1
if category == "Other" and len(other_samples) < 20:
other_samples.append(classify.strip_html(comment)[:90])
other = dist["Other"]
other_pct = other * 100.0 / classified if classified else 0.0
with capsys.disabled():
print(f"\n=== APM classifier smoke test: {FIXTURE.name} ===")
print(
f"rows={len(rows)} classified={classified} "
f"blank-excluded={blank} (blank-comment rows excluded from the total)"
)
print("--- category distribution ---")
for cat, count in dist.most_common():
print(f" {count:4d} {count * 100.0 / classified:5.1f}% {cat}")
print(f"--- Other: {other} ({other_pct:.2f}%) ---")
for sample in other_samples:
print(f" [Other] {sample}")
print(f"--- mismatches flagged: {mismatches} ---")
assert classified > 0
assert other_pct <= MAX_OTHER_PCT, (
f"deterministic Other {other_pct:.2f}% exceeds {MAX_OTHER_PCT}% — "
"the two-axis ladder regressed"
)
# Two-axis model should comfortably beat the legacy ~17% Other.
assert other_pct < 17.0
@pytest.mark.skipif(not FIXTURE.exists(), reason=f"sample export not found: {FIXTURE}")
def test_known_fixture_rows_in_export():
"""Anchor on real rows found in the canonical export."""
import openpyxl
wb = openpyxl.load_workbook(FIXTURE, read_only=True, data_only=True)
rows = list(wb.active.iter_rows(values_only=True))[1:]
by_status = collections.defaultdict(list)
schedule_confirmed_row = None
report_completion_mismatch = None
for row in rows:
comment = row[COL_LAST_COMMENT]
if comment is None or str(comment).strip() == "":
continue
text = classify.strip_html(comment)
by_status[row[COL_WO_STATUS]].append(row)
if (
schedule_confirmed_row is None
and "schedule confirmed with vendor" in text.lower()
and not (row[COL_HOLD_REASON] or "").strip()
):
schedule_confirmed_row = row
if (
report_completion_mismatch is None
and (row[COL_HOLD_REASON] or "").strip().upper() == "REPORT"
and "performed task" in text.lower()
):
report_completion_mismatch = row
# A "WO schedule confirmed with vendor" row → Schedule Confirmed.
assert schedule_confirmed_row is not None, "fixture lacks a schedule-confirmed row"
cat, _ = classify.classify(
schedule_confirmed_row[COL_WO_STATUS],
schedule_confirmed_row[COL_HOLD_REASON],
schedule_confirmed_row[COL_LAST_COMMENT],
)
assert cat == "Schedule Confirmed"
# Any RCAN row → Cancelled.
assert "RCAN" in by_status, "fixture lacks an RCAN row"
rcan = by_status["RCAN"][0]
cat, _ = classify.classify(
rcan[COL_WO_STATUS], rcan[COL_HOLD_REASON], rcan[COL_LAST_COMMENT]
)
assert cat == "Cancelled"
# A REPORT-hold row whose comment claims completion → mismatch non-None.
if report_completion_mismatch is not None:
_, mm = classify.classify(
report_completion_mismatch[COL_WO_STATUS],
report_completion_mismatch[COL_HOLD_REASON],
report_completion_mismatch[COL_LAST_COMMENT],
)
assert mm is not None