"""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 the two-axis classifier, derives ``is_escalation`` / ``is_action`` / ``mismatch``, writes a per-WO snapshot to ``analytics/dt=YYYY-MM-DD/`` as Parquet, and emits a small ``summary.json`` for the slack-post Lambda to read cheaply. The ``apm_wo_snapshots`` Glue table is CDK-defined with partition projection, so the write needs no Glue catalog access. 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 ANALYTICS_PREFIX = "analytics" # Parquet snapshots — the Glue/Athena table reads this META_PREFIX = "meta" # summary.json/details.json — kept OUT of the table's prefix # 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") _lambda = boto3.client("lambda") 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(), } # Comment snippet length for the modal rows — Slack section text caps at 3000 # chars; keep rows short so a full modal stays well under the block limit. _SNIPPET_LEN = 280 def _build_details(df: pd.DataFrame) -> list[dict]: """Per-WO index for the Slack drill-down modals (read from S3 on demand).""" return [ { "wo_number": r.wo_number, "wo_description": r.wo_description, "site": r.site, "department": r.department, "category": r.category, "last_comment": r.last_comment[:_SNIPPET_LEN], "is_escalation": bool(r.is_escalation), "is_action": bool(r.is_action), "mismatch": r.mismatch, } for r in df.itertuples() ] def _event_dt(record: dict) -> str: """Partition date from the S3 event time, not the Lambda wall-clock — stable across retries and across a midnight boundary (a late-night upload retried after midnight keeps the upload day's partition).""" ts = record.get("eventTime") # ISO-8601, e.g. "2026-05-28T22:23:40.123Z" return ts[:10] if ts else datetime.now(timezone.utc).strftime("%Y-%m-%d") def _process_object(bucket: str, key: str, dt: str) -> dict | None: """Classify one export into the dt partition + write its summary/details. Returns the summary dict, or None if the object isn't a usable export.""" if not key.startswith("raw/") or not key.lower().endswith((".xlsx", ".csv")): print(f"Skipping non-export object: s3://{bucket}/{key}") return None 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(f"No classifiable rows in {key} (all comments blank); nothing written.") return None df["dt"] = dt # Pure Parquet write — no Glue registration. The apm_wo_snapshots table is # CDK-defined with partition projection (Phase 3), so Athena derives the dt # partition from the path and the classifier needs no Glue catalog access. # overwrite_partitions keeps a same-day re-upload idempotent (replaces dt=). wr.s3.to_parquet( df=df, path=f"s3://{bucket}/{ANALYTICS_PREFIX}/", dataset=True, partition_cols=["dt"], mode="overwrite_partitions", ) # summary.json/details.json go under meta/ — NOT analytics/. Athena reads # every object in the table's prefix as Parquet, so JSON there breaks queries. summary = _build_summary(df, dt, key, blank) _s3.put_object( Bucket=bucket, Key=f"{META_PREFIX}/dt={dt}/summary.json", Body=json.dumps(summary, indent=2).encode("utf-8"), ContentType="application/json", ) # details.json — per-WO index the slack-post Lambda reads for drill-down modals. _s3.put_object( Bucket=bucket, Key=f"{META_PREFIX}/dt={dt}/details.json", Body=json.dumps(_build_details(df)).encode("utf-8"), ContentType="application/json", ) print( f"Wrote {len(df)} rows, {summary['escalation_total']} escalations " f"({summary['third_escalation_count']} 3rd), {len(summary['mismatches'])} mismatches." ) return summary def handler(event, context): """Classify EVERY export in the S3 event (S3 can batch multiple records), then trigger the Slack post once per affected day. Raises on any failure so the event is retried / lands in the DLQ rather than being silently dropped.""" processed: list[dict] = [] dts: set[str] = set() for record in event.get("Records", []): bucket = record["s3"]["bucket"]["name"] key = unquote_plus(record["s3"]["object"]["key"]) dt = _event_dt(record) summary = _process_object(bucket, key, dt) if summary is not None: processed.append( {"key": key, "dt": dt, "classified": summary["classified_total"]} ) dts.add(dt) for dt in sorted(dts): _invoke_slack_post(dt) return {"processed": processed} def _invoke_slack_post(dt: str) -> None: """Async-invoke the slack-post Lambda (name from env), if wired. Best-effort: a Slack failure must not fail the classification that already landed in S3.""" fn = os.environ.get("SLACK_POST_FUNCTION_NAME") if not fn: return try: _lambda.invoke( FunctionName=fn, InvocationType="Event", Payload=json.dumps({"dt": dt}).encode("utf-8"), ) print(f"Invoked slack-post {fn} for dt={dt}") except Exception as exc: # noqa: BLE001 — never let Slack break the pipeline print(f"slack-post invoke failed (non-fatal): {exc}")