apm-wo-analysis/lambdas/classifier/handler.py

304 lines
11 KiB
Python
Raw Normal View History

"""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")
Add Slack post + interactions Lambdas with drill-down modals (Phase 4) Two push surfaces (no App Home) + interactive drill-down, per CLAUDE.md. Block Kit (blockkit.py, pure/offline): build_daily_summary (header, vs-yesterday deltas, escalation breakdown with 3rd highlighted, action/routine, top sites, mismatch callout, category drill buttons + 📊 Open dashboard link, footer), build_escalation_alert (one @here, returns None on zero-3rd — suppression), and build_wo_modal (views.open payload, capped under Slack's 100-block limit). Lambdas: slack_post/handler.py (classifier-invoked: read today/yesterday summary.json, post daily summary, conditionally post the batched alert from details.json) and slack_post/interactions.py (API Gateway: verify Slack signature, filter details.json, views.open the WO modal within the 3s trigger_id window). slackio.py centralizes Secrets Manager creds, the SSM dashboard URL, signature verification, and analytics/ reads — keeping blockkit pure. Classifier: emit analytics/dt=*/details.json (per-WO index for the modals) and async-invoke slack-post after the snapshot write (best-effort; a Slack failure never fails classification). CDK: slack-post + interactions Lambdas (Docker-bundled slack_sdk), HTTP API on apm-wo.seahaven.com (wildcard ACM cert + Route53 alias; signature-verified, so the route is unauthenticated by design), SSM /apm-wo-analysis/grafana-base-url, and scoped IAM (read analytics/, read the Slack secret + dashboard param; classifier granted lambda:InvokeFunction on slack-post). Slack creds live in one Secrets Manager secret apm-wo-analysis/slack-credentials {botToken, signingSecret, channelId}; cdk.json gains cert/zone/domain context. WO drill-downs link to Grafana only — no APM deep-links (per decision). Deliverables for test time: slack/manifest.yaml (app manifest, interactivity request_url = apm-wo.seahaven.com). Tests: tests/test_blockkit.py (30 offline cases — deltas, zero-3rd None, <100 blocks under large inputs, modal truncation/overflow, dashboard URL) and Phase 4 assertions in test_pipeline_synth.py (both Lambdas, the API route/domain/alias, and no broad/write IAM on the Slack roles). 49/49 tests pass; cdk synth green.
2026-05-28 17:48:51 -04:00
_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(),
}
Add Slack post + interactions Lambdas with drill-down modals (Phase 4) Two push surfaces (no App Home) + interactive drill-down, per CLAUDE.md. Block Kit (blockkit.py, pure/offline): build_daily_summary (header, vs-yesterday deltas, escalation breakdown with 3rd highlighted, action/routine, top sites, mismatch callout, category drill buttons + 📊 Open dashboard link, footer), build_escalation_alert (one @here, returns None on zero-3rd — suppression), and build_wo_modal (views.open payload, capped under Slack's 100-block limit). Lambdas: slack_post/handler.py (classifier-invoked: read today/yesterday summary.json, post daily summary, conditionally post the batched alert from details.json) and slack_post/interactions.py (API Gateway: verify Slack signature, filter details.json, views.open the WO modal within the 3s trigger_id window). slackio.py centralizes Secrets Manager creds, the SSM dashboard URL, signature verification, and analytics/ reads — keeping blockkit pure. Classifier: emit analytics/dt=*/details.json (per-WO index for the modals) and async-invoke slack-post after the snapshot write (best-effort; a Slack failure never fails classification). CDK: slack-post + interactions Lambdas (Docker-bundled slack_sdk), HTTP API on apm-wo.seahaven.com (wildcard ACM cert + Route53 alias; signature-verified, so the route is unauthenticated by design), SSM /apm-wo-analysis/grafana-base-url, and scoped IAM (read analytics/, read the Slack secret + dashboard param; classifier granted lambda:InvokeFunction on slack-post). Slack creds live in one Secrets Manager secret apm-wo-analysis/slack-credentials {botToken, signingSecret, channelId}; cdk.json gains cert/zone/domain context. WO drill-downs link to Grafana only — no APM deep-links (per decision). Deliverables for test time: slack/manifest.yaml (app manifest, interactivity request_url = apm-wo.seahaven.com). Tests: tests/test_blockkit.py (30 offline cases — deltas, zero-3rd None, <100 blocks under large inputs, modal truncation/overflow, dashboard URL) and Phase 4 assertions in test_pipeline_synth.py (both Lambdas, the API route/domain/alias, and no broad/write IAM on the Slack roles). 49/49 tests pass; cdk synth green.
2026-05-28 17:48:51 -04:00
# 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",
)
Add Slack post + interactions Lambdas with drill-down modals (Phase 4) Two push surfaces (no App Home) + interactive drill-down, per CLAUDE.md. Block Kit (blockkit.py, pure/offline): build_daily_summary (header, vs-yesterday deltas, escalation breakdown with 3rd highlighted, action/routine, top sites, mismatch callout, category drill buttons + 📊 Open dashboard link, footer), build_escalation_alert (one @here, returns None on zero-3rd — suppression), and build_wo_modal (views.open payload, capped under Slack's 100-block limit). Lambdas: slack_post/handler.py (classifier-invoked: read today/yesterday summary.json, post daily summary, conditionally post the batched alert from details.json) and slack_post/interactions.py (API Gateway: verify Slack signature, filter details.json, views.open the WO modal within the 3s trigger_id window). slackio.py centralizes Secrets Manager creds, the SSM dashboard URL, signature verification, and analytics/ reads — keeping blockkit pure. Classifier: emit analytics/dt=*/details.json (per-WO index for the modals) and async-invoke slack-post after the snapshot write (best-effort; a Slack failure never fails classification). CDK: slack-post + interactions Lambdas (Docker-bundled slack_sdk), HTTP API on apm-wo.seahaven.com (wildcard ACM cert + Route53 alias; signature-verified, so the route is unauthenticated by design), SSM /apm-wo-analysis/grafana-base-url, and scoped IAM (read analytics/, read the Slack secret + dashboard param; classifier granted lambda:InvokeFunction on slack-post). Slack creds live in one Secrets Manager secret apm-wo-analysis/slack-credentials {botToken, signingSecret, channelId}; cdk.json gains cert/zone/domain context. WO drill-downs link to Grafana only — no APM deep-links (per decision). Deliverables for test time: slack/manifest.yaml (app manifest, interactivity request_url = apm-wo.seahaven.com). Tests: tests/test_blockkit.py (30 offline cases — deltas, zero-3rd None, <100 blocks under large inputs, modal truncation/overflow, dashboard URL) and Phase 4 assertions in test_pipeline_synth.py (both Lambdas, the API route/domain/alias, and no broad/write IAM on the Slack roles). 49/49 tests pass; cdk synth green.
2026-05-28 17:48:51 -04:00
# 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",
Add Slack post + interactions Lambdas with drill-down modals (Phase 4) Two push surfaces (no App Home) + interactive drill-down, per CLAUDE.md. Block Kit (blockkit.py, pure/offline): build_daily_summary (header, vs-yesterday deltas, escalation breakdown with 3rd highlighted, action/routine, top sites, mismatch callout, category drill buttons + 📊 Open dashboard link, footer), build_escalation_alert (one @here, returns None on zero-3rd — suppression), and build_wo_modal (views.open payload, capped under Slack's 100-block limit). Lambdas: slack_post/handler.py (classifier-invoked: read today/yesterday summary.json, post daily summary, conditionally post the batched alert from details.json) and slack_post/interactions.py (API Gateway: verify Slack signature, filter details.json, views.open the WO modal within the 3s trigger_id window). slackio.py centralizes Secrets Manager creds, the SSM dashboard URL, signature verification, and analytics/ reads — keeping blockkit pure. Classifier: emit analytics/dt=*/details.json (per-WO index for the modals) and async-invoke slack-post after the snapshot write (best-effort; a Slack failure never fails classification). CDK: slack-post + interactions Lambdas (Docker-bundled slack_sdk), HTTP API on apm-wo.seahaven.com (wildcard ACM cert + Route53 alias; signature-verified, so the route is unauthenticated by design), SSM /apm-wo-analysis/grafana-base-url, and scoped IAM (read analytics/, read the Slack secret + dashboard param; classifier granted lambda:InvokeFunction on slack-post). Slack creds live in one Secrets Manager secret apm-wo-analysis/slack-credentials {botToken, signingSecret, channelId}; cdk.json gains cert/zone/domain context. WO drill-downs link to Grafana only — no APM deep-links (per decision). Deliverables for test time: slack/manifest.yaml (app manifest, interactivity request_url = apm-wo.seahaven.com). Tests: tests/test_blockkit.py (30 offline cases — deltas, zero-3rd None, <100 blocks under large inputs, modal truncation/overflow, dashboard URL) and Phase 4 assertions in test_pipeline_synth.py (both Lambdas, the API route/domain/alias, and no broad/write IAM on the Slack roles). 49/49 tests pass; cdk synth green.
2026-05-28 17:48:51 -04:00
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}
Add Slack post + interactions Lambdas with drill-down modals (Phase 4) Two push surfaces (no App Home) + interactive drill-down, per CLAUDE.md. Block Kit (blockkit.py, pure/offline): build_daily_summary (header, vs-yesterday deltas, escalation breakdown with 3rd highlighted, action/routine, top sites, mismatch callout, category drill buttons + 📊 Open dashboard link, footer), build_escalation_alert (one @here, returns None on zero-3rd — suppression), and build_wo_modal (views.open payload, capped under Slack's 100-block limit). Lambdas: slack_post/handler.py (classifier-invoked: read today/yesterday summary.json, post daily summary, conditionally post the batched alert from details.json) and slack_post/interactions.py (API Gateway: verify Slack signature, filter details.json, views.open the WO modal within the 3s trigger_id window). slackio.py centralizes Secrets Manager creds, the SSM dashboard URL, signature verification, and analytics/ reads — keeping blockkit pure. Classifier: emit analytics/dt=*/details.json (per-WO index for the modals) and async-invoke slack-post after the snapshot write (best-effort; a Slack failure never fails classification). CDK: slack-post + interactions Lambdas (Docker-bundled slack_sdk), HTTP API on apm-wo.seahaven.com (wildcard ACM cert + Route53 alias; signature-verified, so the route is unauthenticated by design), SSM /apm-wo-analysis/grafana-base-url, and scoped IAM (read analytics/, read the Slack secret + dashboard param; classifier granted lambda:InvokeFunction on slack-post). Slack creds live in one Secrets Manager secret apm-wo-analysis/slack-credentials {botToken, signingSecret, channelId}; cdk.json gains cert/zone/domain context. WO drill-downs link to Grafana only — no APM deep-links (per decision). Deliverables for test time: slack/manifest.yaml (app manifest, interactivity request_url = apm-wo.seahaven.com). Tests: tests/test_blockkit.py (30 offline cases — deltas, zero-3rd None, <100 blocks under large inputs, modal truncation/overflow, dashboard URL) and Phase 4 assertions in test_pipeline_synth.py (both Lambdas, the API route/domain/alias, and no broad/write IAM on the Slack roles). 49/49 tests pass; cdk synth green.
2026-05-28 17:48:51 -04:00
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}")