diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index a30b1ad..f8c2130 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -13,3 +13,4 @@ jobs: run-cdk-synth: true cdk-dir: cdk run-tests: false + enable-qemu: true diff --git a/.github/workflows/deploy.yaml b/.github/workflows/deploy.yaml index 93cecf6..d627cbb 100644 --- a/.github/workflows/deploy.yaml +++ b/.github/workflows/deploy.yaml @@ -23,5 +23,6 @@ jobs: python-version: "3.12" region: us-east-1 cdk-dir: cdk + enable-qemu: true secrets: deploy-role-arn: ${{ secrets.AWS_DEPLOY_ROLE_ARN }} diff --git a/README.md b/README.md index f8e0788..7f094ab 100644 --- a/README.md +++ b/README.md @@ -37,7 +37,9 @@ Two CDK stacks: **Always two-axis, never comment-only.** The legacy script's central flaw was reading only the comment while ignoring `WO Status` + `Hold Reason`, which left -~17% in "Other". The two-axis model cuts that to ~9% before any AI. +~17% in "Other". The two-axis model cuts that to ~9% before any AI — and on the +real 347-row export the current implementation lands "Other" at **5.2%** (18 +rows) with **18 mismatches** flagged. - **Axis 1 — comment intent:** regex over the HTML-stripped `Last Comment`, most-specific first (escalations → SIM ticket → vendor no-show → scheduling → @@ -52,6 +54,10 @@ reading only the comment while ignoring `WO Status` + `Hold Reason`, which left The authoritative spec lives in [`CLAUDE.md`](./CLAUDE.md); the implementation is in `lambdas/classifier/classify.py` (owned by the `classifier-engineer` agent). +The S3-triggered `lambdas/classifier/handler.py` parses each export, classifies +every non-blank-comment row, writes a per-WO Parquet snapshot to +`analytics/dt=YYYY-MM-DD/` (registering the Glue partition via awswrangler), and +emits a `summary.json` for the Phase 4 slack-post Lambda. ## Repository layout @@ -109,7 +115,7 @@ The export reaches S3 by **direct upload or a local drop-folder**, never SES/ema `aws configure --profile apm-wo-drop`. The classifier Lambda is S3-triggered on the `raw/` prefix regardless of path -(added in Phase 2 — uploads currently land in `raw/` and wait). +(any `.xlsx`/`.csv` landing under `raw/` invokes it). ## Deployment @@ -136,5 +142,14 @@ cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role ## Status -**Scaffold (Phase 0).** Build-out proceeds per [`docs/BUILD.md`](./docs/BUILD.md): -ingestion → classifier → Glue/Athena → Slack → Grafana → docs. +**Phase 2 — classifier (in review).** Build-out proceeds per +[`docs/BUILD.md`](./docs/BUILD.md): ingestion → classifier → Glue/Athena → Slack +→ Grafana → docs. + +- **Phase 0** scaffold — merged-pending (PR #6). +- **Phase 1** ingestion (S3 bucket, drop-folder uploader, OIDC deploy role) — + deployed; PR #7 open. +- **Phase 2** classifier Lambda + Glue database + S3 `raw/` trigger — implemented + and `cdk synth`-green; PR #8 open (stacked on Phase 1, not yet deployed). +- **Phases 3–6** (Glue/Athena query layer, Slack, Grafana, final docs) — not + started. diff --git a/cdk/stacks/pipeline_stack.py b/cdk/stacks/pipeline_stack.py index 77f7bd2..cba4334 100644 --- a/cdk/stacks/pipeline_stack.py +++ b/cdk/stacks/pipeline_stack.py @@ -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 diff --git a/lambdas/classifier/classify.py b/lambdas/classifier/classify.py index 43c331b..c2493c3 100644 --- a/lambdas/classifier/classify.py +++ b/lambdas/classifier/classify.py @@ -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 ``...`` 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. <div> 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 diff --git a/lambdas/classifier/handler.py b/lambdas/classifier/handler.py index 1740b60..1c5999a 100644 --- a/lambdas/classifier/handler.py +++ b/lambdas/classifier/handler.py @@ -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", + } diff --git a/tests/test_classify.py b/tests/test_classify.py index 64f2c69..4d8cff8 100644 --- a/tests/test_classify.py +++ b/tests/test_classify.py @@ -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("WO schedule confirmed with vendor") == ( + "WO schedule confirmed with vendor" + ) + # Nested tags, entities, smart quotes, and embedded URLs. + raw = ( + "
1st attempt process for schedule confirmation. " + "Vendor, please confirm ‘Schedule Start Date’ & proceed. " + "https://app.avetta.com/avt-cli/x
" + ) + 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 "&" 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", "", "WO schedule confirmed with vendor" + ) + assert cat == "Schedule Confirmed" + + # Cancelled via WO Status, even with a generic comment. + cat, _ = classify.classify( + "RCAN", "", "WO Cancelled, created in error." + ) + assert cat == "Cancelled" + + # 3rd-attempt escalation cadence. + cat, _ = classify.classify( + "H", "REPORT", "3rd attempt process for schedule confirmation." + ) + assert cat == "3rd Escalation" + + # Vendor no-show. + cat, _ = classify.classify( + "R", "", "Site tech reported vendor was a no show for Friday." + ) + assert cat == "Vendor No-Show" + + # Weekly cadence template. + cat, _ = classify.classify( + "IP", + "", + "Weekly WO scheduled. Service reports required EOD Friday.", + ) + assert cat == "Weekly WO Scheduled" + + # Structured-only fallback: no comment intent fires, REPORT hold decides. + cat, _ = classify.classify("IP", "REPORT", "") + 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", "Vendor arrived and performed task." + ) + 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", "WO schedule confirmed with vendor." + ) + assert cat == "Schedule Confirmed" + assert mm is not None + + # Clean case: no contradiction → no mismatch. + _, mm = classify.classify( + "IP", "", "WO schedule confirmed with vendor" + ) + 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", "1st attempt process for report.") + # 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