Compare commits

...

24 commits

Author SHA1 Message Date
Adam Moussa
54ef5400b1
Merge pull request #13 from Sea-Haven-Industries/feature/phase-6-docs
Some checks are pending
Deploy / deploy (push) Waiting to run
Phase 6: docs and re-enable CD on merge
2026-05-29 14:03:59 -04:00
Adam Moussa
6ddcd56f6c
Merge pull request #11 from Sea-Haven-Industries/feature/phase-5-grafana
Phase 5: self-hosted Grafana (EC2, ALB, dashboards-as-code)
2026-05-29 13:58:30 -04:00
Adam Moussa
c41836b319
Merge pull request #10 from Sea-Haven-Industries/feature/phase-4-slack
Phase 4: Slack post + interactions Lambdas (drill-down modals)
2026-05-29 13:55:24 -04:00
Adam Moussa
9ee3d84fe0 Enable QEMU for arm64 Lambda bundling in CI/CD 2026-05-29 13:37:01 -04:00
Adam Moussa
0ca738ca8c Enable QEMU for arm64 Lambda bundling in CI/CD 2026-05-29 13:37:00 -04:00
Adam Moussa
c4b9bd4305 Enable QEMU for arm64 Lambda bundling in CI/CD 2026-05-29 13:36:59 -04:00
Adam Moussa
1a231131bd Revert "Temporarily disable CD deploy during stack merges"
This reverts commit 611979e44e.
2026-05-29 13:07:02 -04:00
Adam Moussa
fee0fe9a2c Merge phase-0 to carry the CD-disable into phase-6 for an explicit revert 2026-05-29 13:06:57 -04:00
Adam Moussa
dec02170e6 Mark Phase 6 Confluence/Slack docs complete 2026-05-29 12:59:22 -04:00
Adam Moussa
dcda61e728 Add operational runbook (Phase 6)
docs/RUNBOOK.md: incident runbook for a missing daily analysis (detection →
context → triage → resolution-by-cause → post-incident), plus operational
procedures — export upload (direct + drop-folder agent), Grafana OS/app/plugin
patching cadence (clean-replacement preferred), dashboard-JSON redeploy flow +
gotchas, config + grafana.db/EBS backup-restore (DLM snapshot), and a common-
failures quick index. Mirrors to Confluence.
2026-05-29 12:31:14 -04:00
Adam Moussa
60e0b878e3 Refresh README to full operational doc (Phase 6)
Rewrite to the Sea Haven operational template with real resource names from both
stacks: AWS Resources + Lambda Functions tables, Configuration (Secrets/SSM/env/
context), Operations (verify, logs, classifier DLQ, reprocess, Grafana admin),
Documentation, and Notes/Gotchas incl. the deploy-time lessons + a known-debt
list. Corrects stale bits: ALB is internet-facing + office-IP-restricted (not
internal); classifier uses partition projection (no runtime Glue registration);
summary/details JSON live under meta/ not analytics/; grafana uses an instance
role. Adds the resources missing from the old README (classifier DLQ, slack-
interactions Lambda, HTTP API, meta/ + grafana-config/ prefixes, DLM backup,
encrypted volume).
2026-05-29 12:27:05 -04:00
Adam Moussa
eae4d67e00 Apply cross-review findings (Phase 2/4/5 hardening)
From the cross_reviewer (GPT-4.1) per-PR passes, now that the orchestrator is
back up:

Phase 2 (classifier):
- Process ALL S3 records, not just event["Records"][0] — batched notifications
  no longer silently dropped (the review's only BLOCK).
- Derive the partition dt from the S3 event time, not the Lambda wall-clock —
  stable across retries / the midnight boundary.
- Add an SQS dead-letter queue so a failed run surfaces instead of dropping a
  day's data after Lambda's retries.

Phase 4 (Slack):
- Stage throttling (rate 10 / burst 20) on the public /slack/interactions HTTP
  API. (AWS WAF doesn't attach to apigwv2 HTTP APIs; stage throttling is the
  mechanism.)

Phase 5 (Grafana):
- Explicit encrypted=True on the gp3 root volume.

Tests: synth assertions for the DLQ, stage throttling, and the encrypted volume.
60/60 pass; cdk synth green for both stacks. Deferred NITs (print->logging, sig-
failure source-IP logging, S3 versioning, CIDR-maintenance runbook) -> Phase 6.

NOTE: like the earlier deploy fixes these sit on phase-5 but span phases — the
classifier/DLQ to #8, throttling to #10, encryption to #11 — reconcile at merge.
The encrypted-volume change needs the deferred clean instance replacement to
take effect (can't encrypt a live volume in place).
2026-05-29 11:29:14 -04:00
Adam Moussa
c970635be1 Fix WO table filter: drop custom allValue so :singlequote expands All
The 5 multi-select filter vars had allValue='All'. Grafana does NOT apply the
:singlequote format to a custom allValue, so 'All' was injected bare into
site IN (ALL) -> Athena read ALL as a column ('Column ALL cannot be resolved').
Removing the custom allValue lets :singlequote expand the All selection to the
real quoted value list, so the IN clause is valid SQL.
2026-05-29 11:13:10 -04:00
Adam Moussa
cb01b892bc Fix Grafana panels: rawSQL not rawSql (Athena plugin query key)
All 13 panel/variable queries keyed the SQL as rawSql (lowercase); the
grafana-athena-datasource plugin reads rawSQL (capital SQL). With the wrong key
the plugin saw an empty query, so no Athena query ever fired — variables had no
options and every panel showed a clean 'No data' (no error). This was the root
cause of the empty dashboard; data/datasource/permissions were all fine.
2026-05-29 10:55:02 -04:00
Adam Moussa
d884de8fcc Fix Grafana config-sync: --exact-timestamps for same-size updates
aws s3 sync skips same-size files on download unless --exact-timestamps is set,
so a dashboard edit that doesn't change file size (e.g. refresh 2->1, or a query
tweak) never propagated to the instance. Add --exact-timestamps to all four
sync invocations (boot + 15-min timer).
2026-05-28 19:06:06 -04:00
Adam Moussa
dce53fdc86 Fix Grafana template vars: refresh on dashboard load, not time-range change
All 6 query variables had refresh=2 (on time-range change) with no cached value,
so a plain dashboard load never populated them — $dt resolved to empty and every
panel filtered WHERE dt='' (no data). Set refresh=1 (on dashboard load).
2026-05-28 19:01:24 -04:00
Adam Moussa
3c8b7704f6 Fix Grafana Athena auth: use default credential chain, not ec2_iam_role
Grafana rejected the datasource with 'trying to use non-allowed auth method
ec2_iam_role: Failed to create client' — the plugin's allowed_auth_providers
defaults to default,keys,credentials and excludes ec2_iam_role. Switch authType
to 'default' (AWS SDK default chain), which on EC2 resolves to the instance role
via IMDS (still no static keys) and is allowed out of the box.
2026-05-28 18:59:09 -04:00
Adam Moussa
68f3acded8 Fix dashboard rendering: barchart panel type + metadata off the table prefix
Two issues found loading the deployed dashboard:

1. Panels used type "bar-chart" (hyphenated); Grafana's core panel is "barchart"
   — hence "plugin bar-chart required". Fixed both panels.

2. ALL panels showed "no data" because Athena failed with HIVE_BAD_DATA:
   the classifier wrote summary.json/details.json INTO analytics/dt=*/ — the
   same prefix the Glue table scans — so Athena tried to read the JSON as
   Parquet and every query failed. Move the metadata to a separate meta/dt=*/
   prefix: classifier writes there (grant_read_write meta/*), the Slack Lambdas
   read there (read_meta_json, grant_read meta/*), and analytics/ holds only
   Parquet. Verified: the category GROUP BY query now succeeds against Athena.
2026-05-28 18:49:41 -04:00
Adam Moussa
ea54cb1e60 Fix Grafana bootstrap: grafana-cli --homepath + valid plugin version
Found on first boot (cloud-init errored, grafana-server never started):
- grafana-cli needs --homepath=/usr/share/grafana or it can't find config
  defaults; under set -e that aborted the whole bootstrap.
- pinned plugin version 2.18.2 doesn't exist (conflated with awswrangler's
  version) — grafana-athena-datasource latest is 3.2.0.

Verified by running the corrected bootstrap on the instance via SSM: plugin
installs, grafana-server active, /api/health 200, ALB target healthy.

NOTE: the running instance was repaired in-place (the user-data change updated
the launch template but did not replace the instance). The committed user-data
is now correct, so a fresh launch boots clean — a one-time clean instance
replacement should validate that before prod sign-off.
2026-05-28 18:43:08 -04:00
Adam Moussa
f193754c27 Fix deploy-time failures found in prod testing
Two issues only a real deploy/run surfaced (synth + offline tests passed):

1. Classifier exceeded Lambda's 250 MB unzipped limit (bundled awswrangler +
   pandas + pyarrow + numpy). Move them to the AWS-managed SDK-for-pandas layer
   (AWSSDKPandas-Python312-Arm64:27, awswrangler 3.16.1, pre-stripped to fit);
   bundle only openpyxl. Drop the unused anthropic SDK — _call_haiku uses stdlib
   urllib. Function package now ~890 KB.

2. Slack rejected the daily post with invalid_blocks: every category drill
   button shared action_id "drill_category". Qualify it as "drill_category:<cat>"
   for uniqueness; the interactions handler now matches on the prefix. Add a
   regression test asserting all daily-summary action_ids are unique.

Verified in prod: classifier writes Parquet + summary.json + details.json;
slack-post posts the daily summary + 3rd-escalation alert; the interactions
endpoint (apm-wo.seahaven.com) returns 401 on a bad signature. 58/58 tests pass.

NOTE: these fixes sit on the phase-5 branch but logically belong to earlier
phases — the layer fix to #8 (classifier), the Slack fix to #10 — and must be
moved/cherry-picked there before those PRs merge independently. See cleanup.
2026-05-28 18:29:16 -04:00
Adam Moussa
c04c239774 Update README status for Phase 5 Grafana stack 2026-05-28 18:06:10 -04:00
Adam Moussa
4870784fbd Add self-hosted Grafana stack: EC2, ALB, dashboards-as-code (Phase 5)
The one non-serverless piece — Grafana OSS on a t4g.small (AL2023, ARM64) in the
imported seahaven-vpc, fronted by an internet-facing ALB locked by SG to the
office CIDRs (no Client VPN exists, so "VPN-only" = office-IP restriction, the
syslog-server pattern). Instance in private subnets, reachable only from the ALB
SG, administered via SSM Session Manager (no SSH/key pair).

grafana_stack.py: ALB (HTTPS, *.seahaven.com cert, open=False so the SG office
rules aren't undone by an auto 0.0.0.0/0), instance role (Athena query + Glue
read + S3 analytics/athena-results, no static keys), Route53 grafana.seahaven.com
alias, gp3 root volume RETAINed, daily DLM snapshot of the tagged instance, and a
BucketDeployment that uploads grafana/ to the S3 config prefix.

grafana_userdata.sh: install Grafana OSS, pin the Athena datasource plugin, write
grafana.ini (root_url grafana.seahaven.com, kiosk embedding), sync provisioning +
dashboards from S3 on boot, and a systemd timer re-syncs every 15 min so repo
edits land without an instance rebuild.

Dashboard (grafana-author agent, grafana/dashboards/apm-work-orders.json, uid
apm-wo so the Slack 📊 button resolves): 7 panels — category distribution,
escalation summary, action/routine, escalations-by-site, trend time-series over
dt (the new capability), filterable WO table (5 template vars, escalation row
coloring, CSV export, no APM links), and the mismatch panel. Datasource uid
"athena" pinned in the provisioning yaml.

Tests: tests/test_grafana_synth.py — ALB admits only the office CIDRs on 443
(caught and fixed a default 0.0.0.0/0 listener rule), instance only-from-ALB,
no static keys, scoped instance role + SSM, gp3+retained root volume, daily DLM
backup, grafana.seahaven.com alias. 57/57 tests pass; full cdk synth green.
2026-05-28 18:05:20 -04:00
Adam Moussa
cfb9fb85d2 Update README for Phase 4 Slack surfaces 2026-05-28 17:49:57 -04:00
Adam Moussa
c3a2c936d1 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
18 changed files with 3226 additions and 143 deletions

264
README.md
View file

@ -2,8 +2,10 @@
Daily analysis of Amazon **APM work-order** "Last Comment" data for Sea Haven
facility ops. A curated daily filter-view export (~350 work orders) is classified
on two axes, pushed to Slack, and surfaced in a self-hosted Grafana dashboard.
Replaces a legacy Google Apps Script + versioned-Google-Sheet workflow.
on two axes (comment intent + structured `WO Status`/`Hold Reason`), pushed to
Slack as a daily summary plus a batched 3rd-escalation alert, and surfaced in a
self-hosted Grafana dashboard. Replaces a legacy Google Apps Script +
versioned-Google-Sheet workflow.
This is a **sibling concern** to the `apm@` email pipeline in `procurement-ingest`
— it consumes a different feed (the manual export) and does **not** read those
@ -17,79 +19,123 @@ backed by S3 + Athena (Grafana cannot query DynamoDB).
```
APM export (xlsx/csv)
→ S3 raw/ (direct upload OR local launchd drop-folder)
→ S3 raw/ (direct `aws s3 cp` OR local launchd drop-folder)
→ classifier Lambda (HTML strip + two-axis classify, Haiku fallback)
→ S3 analytics/dt=YYYY-MM-DD/ (per-WO daily snapshot, Parquet)
→ Glue table → Athena → Grafana (self-hosted EC2, VPN-only, kiosk)
→ slack-post Lambda (reads today + yesterday partitions)
├→ S3 analytics/dt=YYYY-MM-DD/ (per-WO snapshot, Parquet)
│ → Glue table (partition projection) → Athena → Grafana (EC2, office-IP, kiosk)
├→ S3 meta/dt=YYYY-MM-DD/ (summary.json + details.json — NOT in the table prefix)
└→ async-invoke slack-post Lambda
→ daily summary post [📊 Open dashboard button]
→ standalone batched 3rd-escalation alert (suppressed if zero)
→ standalone batched 3rd-escalation alert (suppressed when zero)
→ drill-down modals via apm-wo.seahaven.com (signature-verified)
```
Two CDK stacks:
Two CDK stacks (`cdk/app.py` instantiates both):
| Stack | Resources |
|---|---|
| `apm-wo-analysis-pipeline` | S3 exports bucket, classifier + slack-post Lambdas, Glue database, Athena workgroup, IAM |
| `apm-wo-analysis-grafana` | EC2 (Grafana OSS), internal ALB, security group, Route53, Athena datasource IAM role |
| `apm-wo-analysis-pipeline` | S3 exports bucket, drop-uploader IAM user, classifier + slack-post + slack-interactions Lambdas, classifier DLQ, Glue DB + projection table, Athena workgroup, HTTP API (`apm-wo.seahaven.com`), SSM param, scoped IAM |
| `apm-wo-analysis-grafana` | EC2 (Grafana OSS), internet-facing **office-IP-restricted** ALB (`grafana.seahaven.com`), security groups, instance IAM role, Route53 alias, daily DLM snapshot, dashboards-as-code S3 deployment |
## AWS Resources
| Resource | Name | Purpose |
|---|---|---|
| S3 bucket | `apm-wo-analysis-exports-328440206208` | Single bucket. Prefixes: `raw/` (incoming, 90-day expiry), `analytics/` (Parquet snapshots, kept), `meta/` (summary/details JSON), `athena-results/` (query output, 30-day expiry), `grafana-config/` (dashboards-as-code). SSE-S3, BPA-all, enforce-SSL, `RETAIN`. |
| IAM user | `apm-wo-drop-uploader` | Drop-folder identity; `s3:PutObject` on `raw/*` only. Access key created out-of-band, stored in local `apm-wo-drop` profile. |
| Glue database | `apm_wo_analysis` | Analytics catalog. |
| Glue table | `apm_wo_snapshots` | External Parquet table over `analytics/`, **partition projection** on `dt` (date, `2026-01-01..NOW`) — no crawler, no `MSCK`. 17-column snapshot schema. |
| Athena workgroup | `apm-wo-analysis` | Enforced result location `athena-results/`, SSE-S3. |
| SQS queue | `apm-wo-analysis-classifier-dlq` | Dead-letter for failed classifier async invocations (14-day retention). |
| SSM parameter | `/apm-wo-analysis/grafana-base-url` | Grafana dashboard URL for the 📊 button / modal links (ops-editable, no redeploy). |
| HTTP API + domain | `apm-wo.seahaven.com` → `POST /slack/interactions` | Slack interactivity endpoint. Stage throttled 10 rps / 20 burst. `*.seahaven.com` ACM cert; Route53 alias. |
| EC2 instance | Grafana (`t4g.small`, AL2023, ARM64) | Self-hosted Grafana OSS in `seahaven-vpc` private subnets, IMDSv2-only, SSM-managed. gp3 20 GB **encrypted**, `DeleteOnTermination=false`, tagged `apm-grafana-backup`. |
| ALB | Grafana ALB (`grafana.seahaven.com`) | Internet-facing, HTTPS 443, SG admits **only office CIDRs** (`47.21.61.4/32`, `96.250.164.146/32`); forwards to instance:3000, health `/api/health`. |
| DLM policy | Grafana volume backup | Daily snapshot (07:00 UTC) of the tagged instance, 7 retained. |
| Route53 | `apm-wo.seahaven.com`, `grafana.seahaven.com` | Aliases in zone `Z06652411XKH89KTZD3XA` (`seahaven.com`). |
## Lambda Functions
All Python 3.12, ARM64, explicit LogGroup (`/aws/lambda/<name>`, 60-day retention).
| Function | Trigger | Purpose |
|---|---|---|
| `apm-wo-analysis-classifier` | S3 `ObjectCreated` on `raw/*.xlsx|.csv` | Parse export, two-axis classify each non-blank-comment row, write per-WO Parquet to `analytics/dt=…/` and `summary.json`/`details.json` to `meta/dt=…/`, then async-invoke slack-post. 512 MB / 120 s. AWS-managed SDK-for-pandas layer (`AWSSDKPandas-Python312-Arm64:27`); DLQ attached. |
| `apm-wo-analysis-slack-post` | Async-invoked by the classifier (`{"dt": …}`) | Read today's + yesterday's `meta/.../summary.json`, post the daily summary, and (only when `third_escalation_count > 0`) the batched 3rd-escalation alert from `details.json`. 256 MB / 30 s. |
| `apm-wo-analysis-slack-interactions` | HTTP API `POST /slack/interactions` | Verify the Slack request signature, read `meta/.../details.json`, and `views.open` a filtered WO-list modal within Slack's 3 s `trigger_id` window. 256 MB / 30 s. |
## Configuration
### Secrets Manager (names only — created out-of-band, never in CloudFormation)
| Secret | Purpose |
|---|---|
| `apm-wo-analysis/anthropic-api-key` | Claude Haiku fallback for ambiguous free-text comments. |
| `apm-wo-analysis/slack-credentials` | JSON `{ botToken, signingSecret, channelId }` for the reused Slack app. |
### SSM Parameters
| Parameter | Purpose |
|---|---|
| `/apm-wo-analysis/grafana-base-url` | Grafana dashboard deep-link base (`https://grafana.seahaven.com/d/apm-wo/...`). |
### Environment Variables (non-secret)
- **classifier:** `APM_HAIKU_FALLBACK` (`on`/`off`), `SLACK_POST_FUNCTION_NAME`.
- **slack-post / slack-interactions:** `SLACK_SECRET_NAME`, `DASHBOARD_URL_PARAM`, `ANALYTICS_BUCKET`.
### GitHub repo secret
- `AWS_DEPLOY_ROLE_ARN` — the OIDC deploy role `githubdeploy-apm-wo-analysis`.
### CDK context (`cdk/cdk.json`)
`wildcardCertArn`, `hostedZoneId`/`hostedZoneName`, `slackInteractionsDomain`, `grafanaDomain`, `grafanaVpcId`/`grafanaAzs`/`grafana{Public,Private}SubnetIds`, `officeCidrs`, `athenaPluginVersion` (`3.2.0`).
## The classification model
**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 — and on the
real 347-row export the current implementation lands "Other" at **5.2%** (18
rows) with **18 mismatches** flagged.
real 347-row export the 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 →
reports → completion → … → other).
- **Axis 2 — structured state:** `Hold Reason` → category and `WO Status`
- **Axis 2 — structured state:** `Hold Reason` → category, and `WO Status`
signals (`RCAN`→Cancelled, `H` corroborates On Hold, `IP`/`R`/`RR` in-flight).
- **Resolution:** comment intent wins when confident → else structured state →
else `Other`. A **Claude Haiku** fallback (Secrets Manager) is reserved for
else `Other`. A **Claude Haiku** fallback (Secrets Manager key) is reserved for
ambiguous free-text with no structured signal.
- **Mismatch detector (a feature):** flags when comment intent contradicts
structured state. Surfaced, never suppressed.
structured state (e.g. "completed" while `WO Status` is `IP`). Surfaced, never
suppressed.
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.
Authoritative spec: [`CLAUDE.md`](./CLAUDE.md). Implementation:
`lambdas/classifier/classify.py` (logic) and `lambdas/classifier/handler.py`
(S3 → Parquet + JSON + invoke).
## Repository layout
```
cdk/
app.py CDK entry point — instantiates both stacks
cdk.json
cdk.json context: cert, zone, subnets, office CIDRs, plugin version
requirements.txt aws-cdk-lib==2.253.1, constructs>=10.6.0
assets/
grafana_userdata.sh EC2 bootstrap: install Grafana + Athena plugin, S3 config sync
stacks/
pipeline_stack.py S3, Lambdas, Glue, Athena, IAM
grafana_stack.py VPC import, EC2, ALB, SG, Route53, datasource role
pipeline_stack.py S3, Lambdas, DLQ, Glue, Athena, HTTP API, IAM
grafana_stack.py VPC import, EC2, ALB, SG, Route53, instance role, DLM
lambdas/
classifier/ S3-triggered: parse → two-axis classify → Parquet
slack_post/ builds + posts the daily summary and alert
classifier/ S3-triggered: parse → two-axis classify → Parquet + meta JSON
slack_post/ blockkit.py (builders), handler.py (post), interactions.py
(modals), slackio.py (Secrets/SSM/S3/signature)
grafana/
provisioning/ Athena datasource + dashboard provider (as code)
dashboards/ committed dashboard JSON (source of truth)
dashboards/ apm-work-orders.json (uid apm-wo) — source of truth
slack/manifest.yaml Slack app manifest (interactivity request URL)
scripts/ local drop-folder uploader + launchd plist
tests/ classifier smoke test
tests/ classifier smoke test + offline synth/blockkit assertions
docs/BUILD.md phased, end-to-end build guide
```
## Configuration
| Where | What |
|---|---|
| **Secrets Manager** | `apm-wo-analysis/anthropic-api-key` (Haiku fallback). Slack bot token reused from the `payments-dashboard` app. |
| **SSM Parameter Store** | operational config (Slack channel ID, schedule expressions, deep-link base URL). |
| **GitHub repo secret** | `AWS_DEPLOY_ROLE_ARN` — the OIDC deploy role `githubdeploy-apm-wo-analysis`. |
No secrets in Lambda environment variables.
## Ingestion (no email)
The export reaches S3 by **direct upload or a local drop-folder**, never SES/email.
@ -97,62 +143,120 @@ The export reaches S3 by **direct upload or a local drop-folder**, never SES/ema
- **Direct:** `aws s3 cp ./export.xlsx s3://apm-wo-analysis-exports-328440206208/raw/`
- **Drop-folder (optional zero-touch):** a launchd agent (`scripts/apm-wo-uploader.sh`
+ `scripts/com.seahaven.apm-wo-uploader.plist`) that watches `~/apm-wo-drop/`,
uploads new `.xlsx`/`.csv` files to `raw/`, and archives them to `uploaded/`.
It uploads with the scoped `apm-wo-drop` AWS profile (IAM user
`apm-wo-drop-uploader` — `s3:PutObject` on `raw/*` only).
uploads new `.xlsx`/`.csv` files to `raw/`, and archives them locally. Uploads
with the scoped `apm-wo-drop` profile (IAM user `apm-wo-drop-uploader`).
Install (the runnable copy **must** live outside `~/Documents` — macOS TCC
sandbox; a repo-path script fails silently with `LastExitStatus=32256`):
Install — the runnable copy and watched folder **must** live outside `~/Documents`
(macOS TCC sandbox; a repo-path script fails silently with `LastExitStatus=32256`):
```bash
install -d "$HOME/.local/bin" "$HOME/apm-wo-drop"
cp scripts/apm-wo-uploader.sh "$HOME/.local/bin/apm-wo-uploader.sh"
chmod +x "$HOME/.local/bin/apm-wo-uploader.sh"
cp scripts/com.seahaven.apm-wo-uploader.plist "$HOME/Library/LaunchAgents/"
launchctl load -w "$HOME/Library/LaunchAgents/com.seahaven.apm-wo-uploader.plist"
aws configure --profile apm-wo-drop # one-time, with the uploader's access key
```
Re-copy the script to `~/.local/bin` after editing the repo source. Configure
the profile once with the uploader's access key:
`aws configure --profile apm-wo-drop`.
The classifier Lambda is S3-triggered on the `raw/` prefix regardless of path
(any `.xlsx`/`.csv` landing under `raw/` invokes it).
## Deployment
CI/CD via the org reusable workflows (no manual prod deploys):
- **CI** (`.github/workflows/ci.yaml`) → `ci-python-sam.yaml@main`: ruff + `cdk synth`.
- **Deploy** (`.github/workflows/deploy.yaml`) → `cd-cdk.yaml@main`: OIDC assume-role,
`cdk deploy --all`, single-flight concurrency.
The OIDC deploy role must exist **before** the first deploy. Deploy order:
```bash
cd cdk && pip install -r requirements.txt
cdk deploy apm-wo-analysis-pipeline # S3, Glue, Athena, Lambdas, IAM
cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role
```
Any `.xlsx`/`.csv` landing under `raw/` invokes the classifier.
## Local development
- pyenv Python 3.12; `ruff check` + `ruff format --check` before pushing (hook-enforced).
- pyenv Python 3.12. `ruff check` + `ruff format --check` before pushing (hook-enforced).
- Tests: `python -m pytest tests/ -q` (60 tests — classifier smoke test against a real
export + offline `cdk.assertions` synth checks + Block Kit builders). No AWS needed.
- Smoke-test the classifier against a **real export** before declaring any
classification change done: `~/Downloads/_documents/Sheet1-1.xlsx`.
- `cdk synth` must pass in CI before merge.
- `cdk synth` must pass in CI before merge (**Docker required** — Lambda deps are
bundled for ARM64).
## Deployment
CI/CD via the org reusable workflows (no manual prod deploys in steady state):
- **CI** (`.github/workflows/ci.yaml`) → `ci-python-sam.yaml@main`: ruff + `cdk synth`. Runs on PRs into `main`.
- **Deploy** (`.github/workflows/deploy.yaml`) → `cd-cdk.yaml@main`: OIDC assume-role, `cdk deploy --all`, single-flight concurrency. Runs on push to `main`.
Stack name/region/account: `apm-wo-analysis-{pipeline,grafana}` / us-east-1 / 328440206208.
OIDC deploy role `githubdeploy-apm-wo-analysis` must exist before the first deploy.
> **Note:** CD is **temporarily disabled** (deploy job gated `if: ${{ false }}` on
> the phase-0 branch) while the Phase 0–5 stack is merged into `main`, to avoid a
> deploy on every merge. **Re-enable as the first Phase 6 step** by reverting that
> commit. Both stacks are already deployed manually and validated in prod.
Manual deploy (emergency/reference; pipeline first so the bucket/table exist):
```bash
cd cdk && pip install -r requirements.txt
cdk deploy apm-wo-analysis-pipeline # S3, Glue, Athena, Lambdas, DLQ, API, IAM
cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, DLM, dashboards
```
## Operations
**Verify it's working** — drop a real export into `raw/`, then within ~1 min:
- `analytics/dt=<today>/*.parquet` and `meta/dt=<today>/{summary,details}.json` appear,
- the daily summary posts to the WO Slack channel (+ alert if any 3rd escalations),
- the Grafana dashboard renders over the office network (`grafana.seahaven.com`).
**Logs:** CloudWatch `/aws/lambda/apm-wo-analysis-{classifier,slack-post,slack-interactions}` (60-day retention).
**Failure modes:**
- Classifier failure (malformed export, transient error) → after Lambda retries, the event lands in **`apm-wo-analysis-classifier-dlq`**. Check the DLQ if a day's data is missing.
- Slack post failure is best-effort and does **not** fail classification (data still lands in S3).
- Grafana down → check the instance via **SSM Session Manager** (no SSH); `systemctl status grafana-server`; ALB target health.
**Reprocess a day:** re-upload the same export to `raw/` — the classifier uses
`overwrite_partitions`, so a same-day re-run replaces that `dt` partition idempotently.
**Grafana admin:** access is office-IP-restricted at the ALB; the instance is
SSM-only. Dashboards are provisioned from `grafana-config/` in S3 (synced on boot
and by a 15-min systemd timer); **edit dashboards in-repo, not in the UI**
(`allowUiUpdates: false`). `grafana.db` lives on the RETAIN'd gp3 volume and is
snapshotted daily by DLM.
## Documentation
- **Confluence "AWS Architecture Map"** (IT space, page **1540098**): a Mermaid
subgraph for this stack — **done** (added under *Serverless Applications*, plus
rows in CI/CD Pipelines, EC2 Inventory, and Key Data Stores).
- **Confluence stack page "APM WO Analysis (apm-wo-analysis)"** (IT space, page
**8749057**, under *AWS Cloud Infrastructure*): full resource/Lambda/secrets
detail — **done**.
- **Slack Apps Inventory** (Confluence page 524569): "APM Work Orders" (new dedicated
app, App ID `A0B6P28V64B`) added — **done**.
- This README + `docs/BUILD.md` (phased build guide) + `CLAUDE.md` (domain spec).
## Notes / Gotchas
- **Partition date** comes from the **S3 event time**, not the Lambda wall-clock —
stable across retries and the midnight boundary.
- **`meta/` vs `analytics/`:** summary/details JSON must stay **out** of `analytics/` —
Athena reads every object in the table prefix as Parquet and chokes on JSON.
- **Classifier deps** (awswrangler/pandas/pyarrow/numpy) come from the AWS-managed
SDK-for-pandas **layer** — bundling them blows Lambda's 250 MB unzipped limit.
Verify the pinned layer ARN/version on region or runtime changes.
- **Grafana dashboard contract:** dashboard `uid` must stay `apm-wo` (the Slack 📊
button deep-links to `d/apm-wo`). Athena plugin query key is `rawSQL` (capital).
Datasource `authType: default` (the EC2 instance role; `ec2_iam_role` is rejected
by the plugin). Config sync uses `aws s3 sync --exact-timestamps` (plain sync
skips same-size edits). Template vars use `refresh: 1` (on load).
- **No Client VPN exists** — "VPN-only" Grafana is realized as **office-IP SG
restriction**. `seahaven-vpc` has a single NAT (one AZ) for instance egress.
- **Slack interactions endpoint** is unauthenticated at the gateway **by design**;
the Lambda verifies the Slack signature (replay window + HMAC). Stage-throttled.
### Known operational debt
- Re-enable CD (revert the phase-0 disable) once the stack is merged.
- One **clean instance replacement** is owed to validate the committed user-data
from a cold boot and to apply root-volume encryption (can't encrypt in place).
- Deferred review NITs: `print()`→`logging`, source-IP logging on signature
failure, S3 versioning, the `'${site:raw}'` WO-table SQL tidy, and a
CIDR-maintenance note in the runbook.
## Status
**Phase 3 — analytics dataset (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).
- **Phase 3** `apm_wo_snapshots` projection table + Athena workgroup —
implemented and `cdk synth`-green; PR #9 open (stacked on Phase 2). The
classifier writes pure Parquet and holds no Glue access (projection handles
partitions).
- **Phases 4–6** (Slack, Grafana, final docs) — not started.
Phases 2–5 implemented, **deployed to prod and validated end-to-end** (classifier,
Slack post + alert + modal, Grafana dashboard). Stacked PRs **#6→#11** are open and
unmerged; cross-review (#2/#4/#5) and `/security-review` of the two public endpoints
are **cleared**. Phase 6 (this docs pass + Confluence + runbook) is in progress on
`feature/phase-6-docs`.

View file

@ -0,0 +1,94 @@
#!/bin/bash
# Grafana OSS bootstrap for the apm-wo-analysis dashboard host (Amazon Linux 2023,
# ARM64). Idempotent enough to re-run. Config + dashboards are pulled from S3
# (the repo is the source of truth); a systemd timer re-syncs dashboards so panel
# updates ship by re-uploading to S3 — no instance rebuild.
#
# Templated by CDK: __CONFIG_BUCKET__ / __CONFIG_PREFIX__ / __PLUGIN_VERSION__.
set -euxo pipefail
CONFIG_BUCKET="__CONFIG_BUCKET__"
CONFIG_PREFIX="__CONFIG_PREFIX__"
PLUGIN_VERSION="__PLUGIN_VERSION__"
# --- Grafana OSS repo + install ---
cat >/etc/yum.repos.d/grafana.repo <<'REPO'
[grafana]
name=grafana
baseurl=https://rpm.grafana.com
repo_gpgcheck=1
enabled=1
gpgcheck=1
gpgkey=https://rpm.grafana.com/gpg.key
sslverify=1
REPO
dnf install -y grafana
# --- Athena datasource plugin (pinned for reproducibility) ---
# --homepath is required or grafana-cli can't find its config defaults.
grafana-cli --homepath=/usr/share/grafana --pluginsDir=/var/lib/grafana/plugins \
plugins install grafana-athena-datasource "${PLUGIN_VERSION}"
# --- grafana.ini: behind the ALB at grafana.seahaven.com, kiosk-friendly ---
cat >/etc/grafana/grafana.ini <<'INI'
[server]
protocol = http
http_port = 3000
root_url = https://grafana.seahaven.com/
enforce_domain = false
[security]
# Behind an office-IP-restricted ALB; allow embedding for the kiosk wall display.
allow_embedding = true
cookie_secure = true
[users]
default_theme = dark
[analytics]
reporting_enabled = false
check_for_updates = false
INI
# --- sync provisioning + dashboards from S3 (repo is source of truth) ---
sync_config() {
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/provisioning/" /etc/grafana/provisioning/ --delete --exact-timestamps
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/dashboards/" /var/lib/grafana/dashboards/ --delete --exact-timestamps
chown -R grafana:grafana /etc/grafana/provisioning /var/lib/grafana/dashboards
}
mkdir -p /var/lib/grafana/dashboards
sync_config
systemctl daemon-reload
systemctl enable --now grafana-server
# --- systemd timer: re-sync dashboards every 15 min so repo edits land without a rebuild ---
cat >/usr/local/bin/grafana-config-sync.sh <<SYNC
#!/bin/bash
set -euo pipefail
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/provisioning/" /etc/grafana/provisioning/ --delete --exact-timestamps
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/dashboards/" /var/lib/grafana/dashboards/ --delete --exact-timestamps
chown -R grafana:grafana /etc/grafana/provisioning /var/lib/grafana/dashboards
SYNC
chmod +x /usr/local/bin/grafana-config-sync.sh
cat >/etc/systemd/system/grafana-config-sync.service <<'SVC'
[Unit]
Description=Sync apm-wo Grafana config/dashboards from S3
[Service]
Type=oneshot
ExecStart=/usr/local/bin/grafana-config-sync.sh
SVC
cat >/etc/systemd/system/grafana-config-sync.timer <<'TIMER'
[Unit]
Description=Periodic apm-wo Grafana config sync
[Timer]
OnBootSec=5min
OnUnitActiveSec=15min
[Install]
WantedBy=timers.target
TIMER
systemctl daemon-reload
systemctl enable --now grafana-config-sync.timer

View file

@ -1,6 +1,17 @@
{
"app": "python3 app.py",
"context": {
"@aws-cdk/core:bootstrapQualifier": "hnb659fds"
"@aws-cdk/core:bootstrapQualifier": "hnb659fds",
"wildcardCertArn": "arn:aws:acm:us-east-1:328440206208:certificate/a66c0994-90d4-410a-a1d9-5595c2a3fae3",
"hostedZoneId": "Z06652411XKH89KTZD3XA",
"hostedZoneName": "seahaven.com",
"slackInteractionsDomain": "apm-wo.seahaven.com",
"grafanaVpcId": "vpc-0d3d4b67bd0cf8a68",
"grafanaAzs": ["us-east-1a", "us-east-1b"],
"grafanaPublicSubnetIds": ["subnet-0eea820effe1b3ae5", "subnet-0012f5895182c1580"],
"grafanaPrivateSubnetIds": ["subnet-04e38c507e96f1926", "subnet-0a0b4fc6f296dfba5"],
"grafanaDomain": "grafana.seahaven.com",
"officeCidrs": ["47.21.61.4/32", "96.250.164.146/32"],
"athenaPluginVersion": "3.2.0"
}
}

View file

@ -1,22 +1,291 @@
"""Grafana stack: VPC import, EC2, ALB, SG, Route53, Athena datasource role.
"""Grafana stack: imported VPC, EC2 (Grafana OSS), ALB, SG, Route53, IAM, backup.
Scaffold — resources are added in Phase 5 of docs/BUILD.md. The running box is
the only non-serverless piece here (self-hosted Grafana OSS on a t4g.small,
ARM64, VPN-only) and carries an OS/Grafana patching + config-backup obligation.
Dashboards are provisioned as code from grafana/ — the running instance is never
the source of truth.
The only non-serverless piece in this repo (self-hosted Grafana OSS on a
t4g.small, ARM64, AL2023). Reachability: an internet-facing ALB whose security
group only admits the office CIDRs — there is no Client VPN in the account, so
"VPN-only" is realized as office-IP restriction (the same pattern syslog-server
uses). The instance sits in private subnets, reachable only from the ALB SG and
administered via SSM Session Manager (no SSH, no key pair).
Dashboards/datasources are provisioned as code: CDK uploads ``grafana/`` to an S3
config prefix; instance user-data syncs it on boot and a systemd timer re-syncs,
so the running box is never the source of truth. The gp3 root volume is RETAINed
and snapshotted daily by DLM; ``grafana.db`` lives there.
Config (cdk.json context): grafanaVpcId, grafanaAzs, grafana{Public,Private}SubnetIds,
grafanaDomain, officeCidrs, wildcardCertArn, hostedZoneId/Name, athenaPluginVersion.
"""
from aws_cdk import Stack
import os
from aws_cdk import (
CfnTag,
Stack,
Tags,
)
from aws_cdk import (
aws_certificatemanager as acm,
)
from aws_cdk import (
aws_dlm as dlm,
)
from aws_cdk import (
aws_ec2 as ec2,
)
from aws_cdk import (
aws_elasticloadbalancingv2 as elbv2,
)
from aws_cdk import (
aws_elasticloadbalancingv2_targets as elbv2_targets,
)
from aws_cdk import (
aws_iam as iam,
)
from aws_cdk import (
aws_route53 as route53,
)
from aws_cdk import (
aws_route53_targets as route53_targets,
)
from aws_cdk import (
aws_s3 as s3,
)
from aws_cdk import (
aws_s3_deployment as s3deploy,
)
from constructs import Construct
EXPORTS_BUCKET = "apm-wo-analysis-exports-328440206208"
CONFIG_PREFIX = "grafana-config"
BACKUP_TAG = "apm-grafana-backup"
GRAFANA_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "grafana")
USERDATA_PATH = os.path.join(
os.path.dirname(__file__), "..", "assets", "grafana_userdata.sh"
)
class GrafanaStack(Stack):
def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None:
super().__init__(scope, construct_id, **kwargs)
# Phase 5 — EC2 (Amazon Linux 2023, Grafana OSS via user-data),
# internal ALB (HTTPS, *.seahaven.com ACM cert),
# SG ingress from VPN/office CIDRs only,
# Route53 alias grafana.seahaven.com,
# instance role: Athena + Glue + S3 read (no static keys). TODO
ctx = self.node.try_get_context
domain = ctx("grafanaDomain")
# Import the shared seahaven-vpc by explicit attributes (no context
# lookup, so the offline synth test needs no AWS credentials).
vpc = ec2.Vpc.from_vpc_attributes(
self,
"SeahavenVpc",
vpc_id=ctx("grafanaVpcId"),
availability_zones=ctx("grafanaAzs"),
public_subnet_ids=ctx("grafanaPublicSubnetIds"),
private_subnet_ids=ctx("grafanaPrivateSubnetIds"),
)
# ----- Security groups -----
alb_sg = ec2.SecurityGroup(
self,
"AlbSg",
vpc=vpc,
description="apm-wo grafana ALB",
allow_all_outbound=True,
)
for cidr in ctx("officeCidrs"):
alb_sg.add_ingress_rule(
ec2.Peer.ipv4(cidr), ec2.Port.tcp(443), f"HTTPS from office {cidr}"
)
instance_sg = ec2.SecurityGroup(
self,
"InstanceSg",
vpc=vpc,
description="apm-wo grafana instance",
allow_all_outbound=True,
)
instance_sg.add_ingress_rule(
alb_sg, ec2.Port.tcp(3000), "Grafana HTTP from the ALB only"
)
# ----- Instance role: Athena query + Glue read + S3 (no static keys) -----
role = iam.Role(
self,
"GrafanaInstanceRole",
assumed_by=iam.ServicePrincipal("ec2.amazonaws.com"),
managed_policies=[
# Session Manager admin access; no SSH / bastion / key pair.
iam.ManagedPolicy.from_aws_managed_policy_name(
"AmazonSSMManagedInstanceCore"
)
],
)
role.add_to_policy(
iam.PolicyStatement(
sid="AthenaQuery",
actions=[
"athena:StartQueryExecution",
"athena:StopQueryExecution",
"athena:GetQueryExecution",
"athena:GetQueryResults",
"athena:GetWorkGroup",
"athena:ListWorkGroups",
],
resources=[
f"arn:aws:athena:{self.region}:{self.account}:workgroup/apm-wo-analysis"
],
)
)
role.add_to_policy(
iam.PolicyStatement(
sid="GlueReadOnly",
actions=[
"glue:GetDatabase",
"glue:GetDatabases",
"glue:GetTable",
"glue:GetTables",
"glue:GetPartition",
"glue:GetPartitions",
],
resources=[
f"arn:aws:glue:{self.region}:{self.account}:catalog",
f"arn:aws:glue:{self.region}:{self.account}:database/apm_wo_analysis",
f"arn:aws:glue:{self.region}:{self.account}:table/apm_wo_analysis/*",
],
)
)
bucket = s3.Bucket.from_bucket_name(self, "Exports", EXPORTS_BUCKET)
# Read the analytics snapshots; read+write athena-results (query output).
bucket.grant_read(role, "analytics/*")
bucket.grant_read_write(role, "athena-results/*")
# Config sync reads the grafana-config prefix.
bucket.grant_read(role, f"{CONFIG_PREFIX}/*")
# ----- Instance (Grafana OSS via user-data) -----
with open(USERDATA_PATH, encoding="utf-8") as fh:
userdata_script = fh.read()
userdata_script = (
userdata_script.replace("__CONFIG_BUCKET__", EXPORTS_BUCKET)
.replace("__CONFIG_PREFIX__", CONFIG_PREFIX)
.replace("__PLUGIN_VERSION__", ctx("athenaPluginVersion"))
)
user_data = ec2.UserData.for_linux()
user_data.add_commands(userdata_script)
instance = ec2.Instance(
self,
"Grafana",
vpc=vpc,
vpc_subnets=ec2.SubnetSelection(
subnet_type=ec2.SubnetType.PRIVATE_WITH_EGRESS
),
instance_type=ec2.InstanceType("t4g.small"),
machine_image=ec2.MachineImage.latest_amazon_linux2023(
cpu_type=ec2.AmazonLinuxCpuType.ARM_64
),
security_group=instance_sg,
role=role,
user_data=user_data,
require_imdsv2=True,
block_devices=[
ec2.BlockDevice(
device_name="/dev/xvda",
# RETAIN the gp3 root volume (grafana.db lives here), encrypted.
volume=ec2.BlockDeviceVolume.ebs(
20,
volume_type=ec2.EbsDeviceVolumeType.GP3,
delete_on_termination=False,
encrypted=True,
),
)
],
)
Tags.of(instance).add(BACKUP_TAG, "true")
# ----- ALB (internet-facing, office-IP-restricted, HTTPS) -----
cert = acm.Certificate.from_certificate_arn(
self, "WildcardCert", ctx("wildcardCertArn")
)
alb = elbv2.ApplicationLoadBalancer(
self,
"Alb",
vpc=vpc,
internet_facing=True,
security_group=alb_sg,
vpc_subnets=ec2.SubnetSelection(subnet_type=ec2.SubnetType.PUBLIC),
)
listener = alb.add_listener(
"Https",
port=443,
protocol=elbv2.ApplicationProtocol.HTTPS,
certificates=[cert],
# Do NOT auto-open 0.0.0.0/0 on 443 — the ALB SG already scopes
# ingress to the office CIDRs. open=True would undo that.
open=False,
)
listener.add_targets(
"GrafanaTarget",
port=3000,
protocol=elbv2.ApplicationProtocol.HTTP,
targets=[elbv2_targets.InstanceTarget(instance, 3000)],
health_check=elbv2.HealthCheck(
path="/api/health", healthy_http_codes="200"
),
)
# ----- Route53 alias grafana.seahaven.com -> ALB -----
zone = route53.HostedZone.from_hosted_zone_attributes(
self,
"SeahavenZone",
hosted_zone_id=ctx("hostedZoneId"),
zone_name=ctx("hostedZoneName"),
)
route53.ARecord(
self,
"GrafanaAlias",
zone=zone,
record_name=domain.split(".")[0],
target=route53.RecordTarget.from_alias(
route53_targets.LoadBalancerTarget(alb)
),
)
# ----- Dashboards-as-code: upload grafana/ to the S3 config prefix -----
s3deploy.BucketDeployment(
self,
"GrafanaConfig",
sources=[s3deploy.Source.asset(GRAFANA_DIR)],
destination_bucket=bucket,
destination_key_prefix=CONFIG_PREFIX,
prune=True,
)
# ----- Daily DLM snapshot of the (tagged) instance's volume -----
dlm_role = iam.Role(
self,
"DlmRole",
assumed_by=iam.ServicePrincipal("dlm.amazonaws.com"),
managed_policies=[
iam.ManagedPolicy.from_aws_managed_policy_name(
"service-role/AWSDataLifecycleManagerServiceRole"
)
],
)
dlm.CfnLifecyclePolicy(
self,
"GrafanaBackup",
description="Daily snapshot of the apm-wo grafana volume",
state="ENABLED",
execution_role_arn=dlm_role.role_arn,
policy_details=dlm.CfnLifecyclePolicy.PolicyDetailsProperty(
resource_types=["INSTANCE"],
target_tags=[CfnTag(key=BACKUP_TAG, value="true")],
schedules=[
dlm.CfnLifecyclePolicy.ScheduleProperty(
name="daily",
create_rule=dlm.CfnLifecyclePolicy.CreateRuleProperty(
interval=24, interval_unit="HOURS", times=["07:00"]
),
retain_rule=dlm.CfnLifecyclePolicy.RetainRuleProperty(count=7),
)
],
),
)

View file

@ -2,8 +2,8 @@
Phase 1 — exports bucket + drop-folder uploader IAM user.
Phase 2 — Glue database + S3-triggered classifier Lambda.
Phase 3 — partition-projection Glue table + Athena workgroup (this file).
Phase 4 — slack-post Lambda + scoped IAM (TODO).
Phase 3 — partition-projection Glue table + Athena workgroup.
Phase 4 — slack-post + interactions Lambdas, API Gateway, scoped IAM (this file).
Lambdas: Python 3.12, ARM64, explicit LogGroup with 60-day retention.
No DynamoDB — this is an S3 + Athena analytics workload (see CLAUDE.md).
@ -17,9 +17,18 @@ from aws_cdk import (
RemovalPolicy,
Stack,
)
from aws_cdk import (
aws_apigatewayv2 as apigwv2,
)
from aws_cdk import (
aws_apigatewayv2_integrations as apigwv2_integrations,
)
from aws_cdk import (
aws_athena as athena,
)
from aws_cdk import (
aws_certificatemanager as acm,
)
from aws_cdk import (
aws_glue as glue,
)
@ -32,6 +41,12 @@ from aws_cdk import (
from aws_cdk import (
aws_logs as logs,
)
from aws_cdk import (
aws_route53 as route53,
)
from aws_cdk import (
aws_route53_targets as route53_targets,
)
from aws_cdk import (
aws_s3 as s3,
)
@ -41,12 +56,28 @@ from aws_cdk import (
from aws_cdk import (
aws_secretsmanager as secretsmanager,
)
from aws_cdk import (
aws_sqs as sqs,
)
from aws_cdk import (
aws_ssm as ssm,
)
from constructs import Construct
GLUE_DATABASE = "apm_wo_analysis"
GLUE_TABLE = "apm_wo_snapshots"
ATHENA_WORKGROUP = "apm-wo-analysis"
# AWS-managed SDK-for-pandas layer (awswrangler 3.16.1, py3.12, arm64). Provides
# awswrangler/pandas/pyarrow/numpy pre-stripped to fit the Lambda size limit.
AWSSDKPANDAS_LAYER_ARN = (
"arn:aws:lambda:us-east-1:336392948345:layer:AWSSDKPandas-Python312-Arm64:27"
)
ANTHROPIC_SECRET = "apm-wo-analysis/anthropic-api-key"
SLACK_SECRET = "apm-wo-analysis/slack-credentials"
GRAFANA_URL_PARAM = "/apm-wo-analysis/grafana-base-url"
GRAFANA_DASHBOARD_URL = (
"https://grafana.seahaven.com/d/apm-wo/apm-work-orders?from=now-30d&to=now"
)
LAMBDAS_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "lambdas")
# Snapshot schema — mirrors the per-WO record written by the classifier
@ -197,6 +228,16 @@ class PipelineStack(Stack):
retention=logs.RetentionDays.TWO_MONTHS,
removal_policy=RemovalPolicy.DESTROY,
)
# DLQ for failed async invocations — a daily pipeline must surface a
# failed run (malformed export, transient error) rather than silently
# drop a day's data after Lambda's retries.
classifier_dlq = sqs.Queue(
self,
"ClassifierDlq",
queue_name="apm-wo-analysis-classifier-dlq",
retention_period=Duration.days(14),
enforce_ssl=True,
)
self.classifier_fn = lambda_.Function(
self,
"Classifier",
@ -207,7 +248,17 @@ class PipelineStack(Stack):
memory_size=512,
timeout=Duration.seconds(120),
log_group=classifier_logs,
dead_letter_queue=classifier_dlq,
environment={"APM_HAIKU_FALLBACK": "on"},
# awswrangler/pandas/pyarrow/numpy come from the AWS-managed
# SDK-for-pandas layer (pre-stripped to fit the 250 MB unzipped
# limit, which bundling them ourselves blows). The function package
# only bundles openpyxl; boto3 is in the runtime, urllib is stdlib.
layers=[
lambda_.LayerVersion.from_layer_version_arn(
self, "PandasLayer", AWSSDKPANDAS_LAYER_ARN
)
],
code=lambda_.Code.from_asset(
os.path.join(LAMBDAS_DIR, "classifier"),
bundling=BundlingOptions(
@ -237,8 +288,140 @@ class PipelineStack(Stack):
# Parquet to S3 — it never touches the catalog (Phase 3).
self.exports_bucket.grant_read(self.classifier_fn, "raw/*")
self.exports_bucket.grant_read_write(self.classifier_fn, "analytics/*")
# summary.json/details.json live under meta/ (kept out of the Athena
# table's analytics/ prefix so queries don't read JSON as Parquet).
self.exports_bucket.grant_read_write(self.classifier_fn, "meta/*")
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
# ----- Phase 4 — Slack post + interactions Lambdas -----
# Reused Slack app credentials (botToken/signingSecret/channelId) live in
# one Secrets Manager secret, created out-of-band like the Anthropic key.
slack_secret = secretsmanager.Secret.from_secret_name_v2(
self, "SlackCreds", SLACK_SECRET
)
# Grafana dashboard URL is operational config — SSM so ops can repoint the
# 📊 button / modal links without a redeploy. The d/apm-wo slug is the
# forward contract Phase 5's dashboard must honor.
dashboard_param = ssm.StringParameter(
self,
"GrafanaUrlParam",
parameter_name=GRAFANA_URL_PARAM,
string_value=GRAFANA_DASHBOARD_URL,
)
slack_env = {
"SLACK_SECRET_NAME": SLACK_SECRET,
"DASHBOARD_URL_PARAM": GRAFANA_URL_PARAM,
"ANALYTICS_BUCKET": self.exports_bucket.bucket_name,
}
def _slack_lambda(construct_id: str, fn_name: str, handler_path: str):
"""A slack_post-package Lambda: read analytics/, the Slack secret, and
the dashboard param. slack_sdk is Docker-bundled for ARM64."""
log_group = logs.LogGroup(
self,
f"{construct_id}Logs",
log_group_name=f"/aws/lambda/{fn_name}",
retention=logs.RetentionDays.TWO_MONTHS,
removal_policy=RemovalPolicy.DESTROY,
)
fn = lambda_.Function(
self,
construct_id,
function_name=fn_name,
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler=handler_path,
memory_size=256,
timeout=Duration.seconds(30),
log_group=log_group,
environment=slack_env,
code=lambda_.Code.from_asset(
os.path.join(LAMBDAS_DIR, "slack_post"),
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",
],
),
),
)
# Slack Lambdas read only the daily JSON under meta/ (not the Parquet).
self.exports_bucket.grant_read(fn, "meta/*")
slack_secret.grant_read(fn)
dashboard_param.grant_read(fn)
return fn
slack_post_fn = _slack_lambda(
"SlackPost", "apm-wo-analysis-slack-post", "handler.handler"
)
interactions_fn = _slack_lambda(
"SlackInteractions",
"apm-wo-analysis-slack-interactions",
"interactions.handler",
)
# Classifier async-invokes slack-post after writing the snapshot.
self.classifier_fn.add_environment(
"SLACK_POST_FUNCTION_NAME", slack_post_fn.function_name
)
slack_post_fn.grant_invoke(self.classifier_fn)
# HTTP API for Slack interactivity, on apm-wo.seahaven.com (the request URL
# registered in the Slack app manifest). The interactions Lambda verifies
# the Slack signature itself; the route is intentionally unauthenticated.
cert = acm.Certificate.from_certificate_arn(
self, "WildcardCert", self.node.try_get_context("wildcardCertArn")
)
slack_domain_name = self.node.try_get_context("slackInteractionsDomain")
slack_domain = apigwv2.DomainName(
self, "SlackDomain", domain_name=slack_domain_name, certificate=cert
)
slack_api = apigwv2.HttpApi(
self,
"SlackInteractionsApi",
default_domain_mapping=apigwv2.DomainMappingOptions(
domain_name=slack_domain
),
)
slack_api.add_routes(
path="/slack/interactions",
methods=[apigwv2.HttpMethod.POST],
integration=apigwv2_integrations.HttpLambdaIntegration(
"InteractionsIntegration", interactions_fn
),
)
# Stage throttling on the public endpoint — Slack interactions are
# low-volume, so cap rate/burst to blunt abuse/DoS against this
# internet-facing route. (AWS WAF doesn't attach to HTTP APIs; stage
# throttling is the apigwv2 mechanism.) Escape hatch to the default stage.
default_stage = slack_api.default_stage.node.default_child
default_stage.default_route_settings = apigwv2.CfnStage.RouteSettingsProperty(
throttling_rate_limit=10, throttling_burst_limit=20
)
# Route53 alias apm-wo.seahaven.com → the API Gateway custom domain.
zone = route53.HostedZone.from_hosted_zone_attributes(
self,
"SeahavenZone",
hosted_zone_id=self.node.try_get_context("hostedZoneId"),
zone_name=self.node.try_get_context("hostedZoneName"),
)
route53.ARecord(
self,
"SlackDomainAlias",
zone=zone,
record_name=slack_domain_name.split(".")[0],
target=route53.RecordTarget.from_alias(
route53_targets.ApiGatewayv2DomainProperties(
slack_domain.regional_domain_name,
slack_domain.regional_hosted_zone_id,
)
),
)

187
docs/RUNBOOK.md Normal file
View file

@ -0,0 +1,187 @@
# apm-wo-analysis — Operational Runbook
Operational procedures and incident response for the daily APM work-order
analysis pipeline. Account **328440206208** / **us-east-1**. Stacks
`apm-wo-analysis-pipeline` and `apm-wo-analysis-grafana`. See [`README.md`](../README.md)
for architecture and resource detail.
**Admin access:** the Grafana box is **SSM Session Manager only** (no SSH/key pair):
`aws ssm start-session --target <instance-id>`. AWS API via the `office_mac` /
deploy roles. Grafana UI is office-IP-restricted at the ALB.
---
## 1. Incident Runbook: Missing daily APM WO analysis
**Severity: High** — a day's WO analysis is missing: escalations (incl. 3rd-escalation
alerts) aren't surfaced and the dashboard has no new snapshot. Recoverable by
re-processing the export; not permanent data loss.
### Detection
- No **daily summary** in the WO Slack channel by the usual time (or no 3rd-escalation alert on a day one's expected).
- **Grafana** shows no `dt = today` in the Snapshot Date dropdown / panels empty for today.
- **Messages in `apm-wo-analysis-classifier-dlq`** (SQS) — strongest signal the classifier failed.
- CloudWatch errors in `/aws/lambda/apm-wo-analysis-classifier` or `-slack-post`.
### Context
| Item | Value |
|---|---|
| Stacks | `apm-wo-analysis-pipeline`, `apm-wo-analysis-grafana` |
| Lambdas | `apm-wo-analysis-classifier`, `-slack-post`, `-slack-interactions` |
| S3 | `apm-wo-analysis-exports-328440206208` — `raw/`, `analytics/dt=…/`, `meta/dt=…/` |
| SQS DLQ | `apm-wo-analysis-classifier-dlq` |
| Glue / Athena | db `apm_wo_analysis`, table `apm_wo_snapshots`, workgroup `apm-wo-analysis` |
| External | Slack, Anthropic API (Haiku fallback) |
| Secrets (names) | `apm-wo-analysis/slack-credentials`, `apm-wo-analysis/anthropic-api-key` |
| SSM | `/apm-wo-analysis/grafana-base-url` |
### Triage
1. **Export uploaded?** `aws s3 ls s3://apm-wo-analysis-exports-328440206208/raw/` — today's file present? Absent → upstream (§2.1), not the pipeline.
2. **Classifier ran/failed?** `/aws/lambda/apm-wo-analysis-classifier` logs; peek the DLQ: `aws sqs receive-message --queue-url <dlq-url> --max-number-of-messages 1`.
3. **Outputs written?** `aws s3 ls .../analytics/dt=<today>/` (Parquet) and `.../meta/dt=<today>/` (`summary.json`, `details.json`).
4. **slack-post ran/failed?** `/aws/lambda/apm-wo-analysis-slack-post` logs — `SlackApiError` (`invalid_auth`, `not_in_channel`, `invalid_blocks`)?
5. **Slack creds** valid + bot still in channel? (`apm-wo-analysis/slack-credentials`).
6. **IAM/throttle:** grep logs for `AccessDenied` / throttling.
7. **Recent change?** Any merge/deploy to `main` just before the failure.
### Resolution (by root cause)
1. **Export not uploaded** → `aws s3 cp <export>.xlsx s3://apm-wo-analysis-exports-328440206208/raw/`; then check the drop-folder agent (§2.1).
2. **Classifier failed (DLQ)** → read the DLQ message; fix; **reprocess by re-uploading the export to `raw/`** (`overwrite_partitions` makes same-`dt` idempotent).
3. **Outputs present, no Slack post** → re-invoke:
```bash
aws lambda invoke --function-name apm-wo-analysis-slack-post \
--payload "$(printf '{"dt":"<YYYY-MM-DD>"}' | base64)" /tmp/out.json
```
If Slack auth was the cause → rotate `apm-wo-analysis/slack-credentials` (and/or `/invite` the bot), then re-invoke (secret read per-call; no redeploy).
4. **Code regression** → identify the PR, revert/hotfix, redeploy via CI (push to `main`).
5. **Data present, Grafana empty** → datasource/dashboard issue (§2.3); check the instance via SSM.
6. **Anthropic/Haiku down** → non-fatal (deterministic path still classifies ~95%); set classifier env `APM_HAIKU_FALLBACK=off` to bypass.
### Post-Incident
- Verify re-process → summary posts + Grafana shows today's `dt`.
- Check for **other missed days** (gaps in `analytics/dt=…`/`meta/`) and reprocess each.
- Update README/Confluence if knowledge changed; add a memory entry; add a test if code caused it.
---
## 2. Operational Procedures
### 2.1 How the export gets uploaded
The curated daily APM filter-view export (~350 WOs, `.xlsx`/`.csv`) reaches S3 by
**direct upload or a local drop-folder** — never SES/email. Any object under
`raw/` with a `.xlsx`/`.csv` suffix triggers the classifier.
- **Direct:** `aws s3 cp ./export.xlsx s3://apm-wo-analysis-exports-328440206208/raw/`
- **Drop-folder (zero-touch):** a launchd agent (`com.seahaven.apm-wo-uploader`)
watches `~/apm-wo-drop/`, uploads new files to `raw/` using the scoped
`apm-wo-drop` profile (IAM user **`apm-wo-drop-uploader`** — `s3:PutObject` on
`raw/*` only), and archives them locally. The runnable script lives at
`~/.local/bin/apm-wo-uploader.sh` (it and the watched folder **must** be outside
`~/Documents` — macOS TCC sandbox).
**Verify the agent:**
```bash
launchctl list | grep apm-wo-uploader # present + last exit 0
tail -f ~/apm-wo-drop/uploader.log # per-run logging (TBD: confirm log path)
```
**Common upstream issues:** agent unloaded (`launchctl load -w …plist`); script
moved back into `~/Documents` (TCC blocks it — `LastExitStatus=32256`); `apm-wo-drop`
access key expired/rotated (`aws configure --profile apm-wo-drop`).
### 2.2 Grafana OS / app patching cadence
The Grafana EC2 box (`t4g.small`, **Amazon Linux 2023**, ARM64) is the only
patch-bearing piece — everything else is serverless. It's **reproducible from
`cdk/assets/grafana_userdata.sh`**, so the preferred patch path is a **clean
instance replacement** rather than long-lived in-place drift.
- **OS (recommended monthly + on critical CVEs):** via SSM —
`sudo dnf upgrade --security -y && sudo reboot` (Session Manager, or an SSM
Run Command / Patch Manager maintenance window — **TBD: not yet automated**).
- **Grafana OSS:** `sudo dnf upgrade grafana -y && sudo systemctl restart grafana-server`
(installed from the pinned `rpm.grafana.com` repo).
- **Athena datasource plugin:** pinned to **`3.2.0`** in `cdk/cdk.json`
(`athenaPluginVersion`). Bump there, then redeploy/replace the instance.
- **Preferred = clean replacement:** terminate the instance; `cdk deploy
apm-wo-analysis-grafana` relaunches it from the latest AL2023 AMI and re-runs
user-data (fresh Grafana + plugin + config sync). The root volume is
`DeleteOnTermination=false`, so detach/reuse or restore `grafana.db` (§2.4) if
local settings must persist. Validates the committed user-data from a cold boot.
> **Outstanding:** one clean instance replacement is owed to validate cold-boot
> user-data and apply root-volume encryption (encryption can't be added in place).
### 2.3 Dashboard-JSON redeploy
Source of truth is **`grafana/dashboards/apm-work-orders.json`** in this repo
(uid **`apm-wo`**); the running instance is never the source of truth
(`allowUiUpdates: false` — UI edits are reverted on the next sync).
**Flow:** edit JSON in repo → `cdk deploy apm-wo-analysis-grafana` (the
`BucketDeployment` uploads `grafana/` to `s3://…/grafana-config/`) → the instance
syncs S3 → `/var/lib/grafana/dashboards/` (on boot + a **15-min systemd timer**)
→ Grafana's file provider polls every **60 s** and reloads.
**Apply immediately** (skip the timer) via SSM:
```bash
sudo /usr/local/bin/grafana-config-sync.sh # pulls grafana-config/ from S3
# Grafana file provider picks up the dashboard within ~60s
```
**Gotchas:**
- The sync uses `aws s3 sync --exact-timestamps` — required so **same-size edits**
(e.g. a one-char query change) actually propagate.
- **Datasource/provisioning** changes (`grafana/provisioning/*.yaml`) are loaded
at **startup** — after syncing, `sudo systemctl restart grafana-server` (a
dashboard-only change does **not** need a restart).
- Athena query key is **`rawSQL`** (capital); datasource `authType: default`;
template vars `refresh: 1`. (See README "Notes / Gotchas".)
### 2.4 Backup & restore (config + EBS / grafana.db)
Two distinct layers:
**Config (dashboards, datasources, provisioning)** — fully **reproducible from
git** (`grafana/` → S3 `grafana-config/`). *Restore:* `cdk deploy
apm-wo-analysis-grafana` (or `grafana-config-sync.sh` on the box). No snapshot needed.
**Local state (`/var/lib/grafana/grafana.db`)** — Grafana's SQLite (admin user,
any API keys, org prefs). Lives on the **gp3 root volume** (encrypted,
`DeleteOnTermination=false`). Backed up by a **daily DLM snapshot** (07:00 UTC,
7 retained) of the instance (tag `apm-grafana-backup=true`), policy in the
grafana stack.
*Restore from snapshot:*
```bash
# find the latest DLM snapshot
aws ec2 describe-snapshots --owner-ids self \
--filters "Name=tag:aws:dlm:lifecycle-policy-id,Values=*" \
--query 'reverse(sort_by(Snapshots,&StartTime))[0].SnapshotId' --output text
# create a volume from it and attach to a replacement instance, OR mount it and
# copy /var/lib/grafana/grafana.db onto the new instance, then:
sudo systemctl restart grafana-server
```
Because dashboards + datasource are provisioned from code, the only thing the
snapshot uniquely protects is `grafana.db` (admin/login state) — low stakes; a
fresh instance + provisioning recovers everything else.
> **Note:** the Grafana **admin auth model is unsettled (TBD)** — the password was
> reset ad-hoc during build/testing. Decide the intended model (fixed admin
> password in Secrets Manager / SSO / anonymous view-only for the kiosk) and
> document it here.
### 2.5 Common failures (quick index)
| Symptom | Likely cause | Go to |
|---|---|---|
| No daily Slack post / no new dashboard day | export not uploaded, classifier failed (DLQ), or slack-post failed | §1 |
| Dashboard loads but all panels "No data" | datasource/auth, `rawSQL`, template `refresh`, or JSON in `analytics/` prefix | §2.3, README gotchas |
| Grafana unreachable | instance down / ALB unhealthy / office IP changed (`officeCidrs`) | SSM triage; `aws elbv2 describe-target-health` |
| Slack modal click does nothing / error | `apm-wo-analysis-slack-interactions`, API Gateway, or signing-secret mismatch | `/aws/lambda/apm-wo-analysis-slack-interactions` logs |
| Exports never arrive in `raw/` | drop-folder agent unloaded / TCC / expired key | §2.1 |
| Deploy not applying | OIDC role, CloudFormation rollback, Docker bundling | CloudFormation events; CI logs |
---
*Maintained in-repo (`docs/RUNBOOK.md`) and mirrored to Confluence. Update both
when operational knowledge changes.*

View file

@ -1,7 +1,7 @@
{
"uid": "apm-wo",
"title": "APM Work Orders",
"description": "Daily APM work-order analysis — breakdown, escalations, trend, filterable WO table, mismatches. Built in Phase 5 (docs/BUILD.md) by the grafana-author agent.",
"description": "Daily APM work-order analysis — breakdown, escalations, trend, filterable WO table, mismatches. Built in Phase 5 by the grafana-author agent. Athena datasource uid='athena'; data grain is one row per WO per daily snapshot partitioned by dt.",
"tags": ["apm", "work-orders"],
"timezone": "browser",
"schemaVersion": 39,
@ -12,7 +12,794 @@
"to": "now"
},
"templating": {
"list": []
"list": [
{
"name": "dt",
"label": "Snapshot Date",
"description": "Single daily snapshot partition. Defaults to the latest available dt. Most panels filter WHERE dt = '$dt'.",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT dt FROM apm_wo_analysis.apm_wo_snapshots ORDER BY dt DESC",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 0,
"multi": false,
"includeAll": false,
"current": {},
"options": [],
"hide": 0
},
{
"name": "site",
"label": "Site",
"description": "Multi-select. WO table applies: AND ('${site:raw}' = 'All' OR site IN (${site:singlequote}))",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT site FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND site<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 1,
"multi": true,
"includeAll": true,
"current": {},
"options": [],
"hide": 0
},
{
"name": "department",
"label": "Department",
"description": "Multi-select. WO table applies: AND ('${department:raw}' = 'All' OR department IN (${department:singlequote}))",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT department FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND department<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 1,
"multi": true,
"includeAll": true,
"current": {},
"options": [],
"hide": 0
},
{
"name": "category",
"label": "Category",
"description": "Multi-select. WO table applies: AND ('${category:raw}' = 'All' OR category IN (${category:singlequote}))",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT category FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND category<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 1,
"multi": true,
"includeAll": true,
"current": {},
"options": [],
"hide": 0
},
{
"name": "wo_status",
"label": "WO Status",
"description": "Multi-select. WO table applies: AND ('${wo_status:raw}' = 'All' OR wo_status IN (${wo_status:singlequote}))",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT wo_status FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND wo_status<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 1,
"multi": true,
"includeAll": true,
"current": {},
"options": [],
"hide": 0
},
{
"name": "hold_reason",
"label": "Hold Reason",
"description": "Multi-select. WO table applies: AND ('${hold_reason:raw}' = 'All' OR hold_reason IN (${hold_reason:singlequote}))",
"type": "query",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT hold_reason FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND hold_reason<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
},
"refresh": 1,
"sort": 1,
"multi": true,
"includeAll": true,
"current": {},
"options": [],
"hide": 0
}
]
},
"panels": []
"panels": [
{
"id": 1,
"type": "barchart",
"title": "Category Distribution",
"description": "Count of work orders by classification category for the selected snapshot date. Sorted descending by volume. Filter using the template variables above.",
"gridPos": { "x": 0, "y": 0, "w": 14, "h": 9 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT category, COUNT(*) AS wos FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' GROUP BY 1 ORDER BY 2 DESC",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"orientation": "horizontal",
"barRadius": 0,
"groupWidth": 0.7,
"showValue": "always",
"stacking": "none",
"tooltip": { "mode": "single", "sort": "none" },
"legend": { "showLegend": false, "displayMode": "list", "placement": "bottom" }
},
"fieldConfig": {
"defaults": {
"color": { "mode": "palette-classic" },
"custom": {
"axisBorderShow": false,
"axisCenteredZero": false,
"axisColorMode": "text",
"axisLabel": "",
"axisPlacement": "auto",
"fillOpacity": 80,
"gradientMode": "none",
"hideFrom": { "legend": false, "tooltip": false, "viz": false },
"lineWidth": 1,
"scaleDistribution": { "type": "linear" },
"thresholdsStyle": { "mode": "off" }
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "green", "value": null },
{ "color": "red", "value": 80 }
]
}
},
"overrides": []
}
},
{
"id": 2,
"type": "piechart",
"title": "Escalation Summary",
"description": "Distribution across escalation categories (1st, 2nd, 3rd Escalation, SIM Ticket, Other Escalation) for the selected snapshot date. 3rd Escalation is highlighted red.",
"gridPos": { "x": 14, "y": 0, "w": 10, "h": 9 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT category, COUNT(*) AS escalations FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND is_escalation=true GROUP BY 1 ORDER BY 2 DESC",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"pieType": "pie",
"displayLabels": ["name", "value"],
"tooltip": { "mode": "single", "sort": "none" },
"legend": { "showLegend": true, "displayMode": "table", "placement": "right", "values": ["value", "percent"] }
},
"fieldConfig": {
"defaults": {
"color": { "mode": "palette-classic" },
"custom": {
"hideFrom": { "legend": false, "tooltip": false, "viz": false }
},
"mappings": []
},
"overrides": [
{
"matcher": { "id": "byName", "options": "3rd Escalation" },
"properties": [
{ "id": "color", "value": { "fixedColor": "red", "mode": "fixed" } }
]
},
{
"matcher": { "id": "byName", "options": "2nd Escalation" },
"properties": [
{ "id": "color", "value": { "fixedColor": "orange", "mode": "fixed" } }
]
},
{
"matcher": { "id": "byName", "options": "1st Escalation" },
"properties": [
{ "id": "color", "value": { "fixedColor": "yellow", "mode": "fixed" } }
]
},
{
"matcher": { "id": "byName", "options": "SIM Ticket" },
"properties": [
{ "id": "color", "value": { "fixedColor": "purple", "mode": "fixed" } }
]
}
]
}
},
{
"id": 3,
"type": "piechart",
"title": "Action-Needed vs Routine",
"description": "Action-needed WOs require attention (escalations, scheduling, vendor/report waits, status inquiries, vendor no-shows). Routine WOs are on track. Counts are for the selected snapshot date.",
"gridPos": { "x": 0, "y": 9, "w": 8, "h": 8 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT CASE WHEN is_action=true THEN 'Action Needed' ELSE 'Routine' END AS action_type, COUNT(*) AS wos FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' GROUP BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"pieType": "donut",
"displayLabels": ["name", "percent"],
"tooltip": { "mode": "single", "sort": "none" },
"legend": { "showLegend": true, "displayMode": "list", "placement": "bottom", "values": ["value", "percent"] }
},
"fieldConfig": {
"defaults": {
"color": { "mode": "palette-classic" },
"custom": {
"hideFrom": { "legend": false, "tooltip": false, "viz": false }
},
"mappings": []
},
"overrides": [
{
"matcher": { "id": "byName", "options": "Action Needed" },
"properties": [
{ "id": "color", "value": { "fixedColor": "semi-dark-orange", "mode": "fixed" } }
]
},
{
"matcher": { "id": "byName", "options": "Routine" },
"properties": [
{ "id": "color", "value": { "fixedColor": "green", "mode": "fixed" } }
]
}
]
}
},
{
"id": 4,
"type": "barchart",
"title": "Escalations by Site",
"description": "Total escalations per site for the selected snapshot date. Includes all escalation categories.",
"gridPos": { "x": 8, "y": 9, "w": 16, "h": 8 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT site, COUNT(*) AS escalations FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND is_escalation=true AND site<>'' GROUP BY 1 ORDER BY 2 DESC",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"orientation": "vertical",
"barRadius": 0,
"groupWidth": 0.7,
"showValue": "always",
"stacking": "none",
"tooltip": { "mode": "single", "sort": "none" },
"legend": { "showLegend": false, "displayMode": "list", "placement": "bottom" },
"xTickLabelRotation": -45,
"xTickLabelMaxLength": 12
},
"fieldConfig": {
"defaults": {
"color": { "mode": "fixed", "fixedColor": "semi-dark-red" },
"custom": {
"axisBorderShow": false,
"axisCenteredZero": false,
"axisColorMode": "text",
"axisLabel": "Escalations",
"axisPlacement": "auto",
"fillOpacity": 80,
"gradientMode": "none",
"hideFrom": { "legend": false, "tooltip": false, "viz": false },
"lineWidth": 1,
"scaleDistribution": { "type": "linear" },
"thresholdsStyle": { "mode": "off" }
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "green", "value": null }
]
}
},
"overrides": []
}
},
{
"id": 5,
"type": "timeseries",
"title": "Trend Over Time — Escalations & Action-Needed per Day",
"description": "Runs across all partitions (no $dt filter) to show daily escalation and action-needed volume over time. Use the dashboard time range picker to zoom in. This is the capability the legacy Sheet never had.",
"gridPos": { "x": 0, "y": 17, "w": 24, "h": 9 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT date_parse(dt, '%Y-%m-%d') AS time, SUM(CAST(is_escalation AS INTEGER)) AS escalations, SUM(CAST(is_action AS INTEGER)) AS action_needed FROM apm_wo_analysis.apm_wo_snapshots GROUP BY 1 ORDER BY 1",
"format": "timeSeries",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"tooltip": { "mode": "multi", "sort": "desc" },
"legend": { "showLegend": true, "displayMode": "list", "placement": "bottom", "calcs": ["mean", "max", "last"] }
},
"fieldConfig": {
"defaults": {
"color": { "mode": "palette-classic" },
"custom": {
"axisBorderShow": false,
"axisCenteredZero": false,
"axisColorMode": "text",
"axisLabel": "Work Orders",
"axisPlacement": "auto",
"barAlignment": 0,
"drawStyle": "line",
"fillOpacity": 10,
"gradientMode": "none",
"hideFrom": { "legend": false, "tooltip": false, "viz": false },
"insertNulls": false,
"lineInterpolation": "linear",
"lineWidth": 2,
"pointSize": 5,
"scaleDistribution": { "type": "linear" },
"showPoints": "auto",
"spanNulls": false,
"stacking": { "group": "A", "mode": "none" },
"thresholdsStyle": { "mode": "off" }
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "green", "value": null }
]
},
"unit": "short"
},
"overrides": [
{
"matcher": { "id": "byName", "options": "escalations" },
"properties": [
{ "id": "color", "value": { "fixedColor": "semi-dark-red", "mode": "fixed" } },
{ "id": "displayName", "value": "Escalations" }
]
},
{
"matcher": { "id": "byName", "options": "action_needed" },
"properties": [
{ "id": "color", "value": { "fixedColor": "semi-dark-orange", "mode": "fixed" } },
{ "id": "displayName", "value": "Action Needed" }
]
}
]
}
},
{
"id": 6,
"type": "table",
"title": "Work Order Detail Table",
"description": "Filterable, exportable table of all work orders for the selected snapshot date. Multi-value template variables are applied via: AND ('${variable:raw}' = 'All' OR column IN (${variable:singlequote})). Escalation rows are color-coded: 3rd Escalation = red, 2nd Escalation = orange, 1st Escalation = yellow. WO numbers are plain text (no APM deep-link per project decision). Use the Download CSV button (table header menu) to export. last_comment is wrapped for readability.",
"gridPos": { "x": 0, "y": 26, "w": 24, "h": 14 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT wo_number, site, department, category, wo_status, hold_reason, wo_description, last_comment FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND ('${site:raw}' = 'All' OR site IN (${site:singlequote})) AND ('${department:raw}' = 'All' OR department IN (${department:singlequote})) AND ('${category:raw}' = 'All' OR category IN (${category:singlequote})) AND ('${wo_status:raw}' = 'All' OR wo_status IN (${wo_status:singlequote})) AND ('${hold_reason:raw}' = 'All' OR hold_reason IN (${hold_reason:singlequote})) ORDER BY category, site",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"frameIndex": 0,
"showHeader": true,
"sortBy": [],
"footer": {
"show": false,
"reducer": ["sum"],
"fields": "",
"enablePagination": false
}
},
"fieldConfig": {
"defaults": {
"color": { "mode": "thresholds" },
"custom": {
"align": "left",
"cellOptions": { "type": "auto" },
"inspect": false,
"filterable": true,
"minWidth": 80,
"width": 0
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "text", "value": null }
]
}
},
"overrides": [
{
"matcher": { "id": "byName", "options": "wo_number" },
"properties": [
{ "id": "displayName", "value": "WO #" },
{ "id": "custom.width", "value": 100 }
]
},
{
"matcher": { "id": "byName", "options": "site" },
"properties": [
{ "id": "displayName", "value": "Site" },
{ "id": "custom.width", "value": 80 }
]
},
{
"matcher": { "id": "byName", "options": "department" },
"properties": [
{ "id": "displayName", "value": "Dept" },
{ "id": "custom.width", "value": 80 }
]
},
{
"matcher": { "id": "byName", "options": "category" },
"properties": [
{ "id": "displayName", "value": "Category" },
{ "id": "custom.width", "value": 160 },
{
"id": "mappings",
"value": [
{
"type": "value",
"options": {
"3rd Escalation": {
"color": "dark-red",
"index": 0
},
"2nd Escalation": {
"color": "dark-orange",
"index": 1
},
"1st Escalation": {
"color": "dark-yellow",
"index": 2
},
"SIM Ticket": {
"color": "dark-purple",
"index": 3
},
"Other Escalation": {
"color": "orange",
"index": 4
}
}
}
]
},
{
"id": "custom.cellOptions",
"value": { "type": "color-text" }
}
]
},
{
"matcher": { "id": "byName", "options": "wo_status" },
"properties": [
{ "id": "displayName", "value": "Status" },
{ "id": "custom.width", "value": 80 }
]
},
{
"matcher": { "id": "byName", "options": "hold_reason" },
"properties": [
{ "id": "displayName", "value": "Hold Reason" },
{ "id": "custom.width", "value": 120 }
]
},
{
"matcher": { "id": "byName", "options": "wo_description" },
"properties": [
{ "id": "displayName", "value": "Description" },
{ "id": "custom.width", "value": 220 }
]
},
{
"matcher": { "id": "byName", "options": "last_comment" },
"properties": [
{ "id": "displayName", "value": "Last Comment" },
{
"id": "custom.cellOptions",
"value": { "type": "auto", "wrapText": true }
},
{ "id": "custom.width", "value": 420 },
{ "id": "custom.minWidth", "value": 200 }
]
}
]
}
},
{
"id": 7,
"type": "table",
"title": "Mismatch Panel — Comment vs Structured-State Contradictions",
"description": "Work orders where the classified comment intent contradicts the structured WO state (e.g. comment says 'completed' but WO is on a REPORT/VENDOR/SCHEDULING hold, or comment says 'scheduled' while status is IP). Non-empty mismatch column only. Use these rows for manual review before the next export.",
"gridPos": { "x": 0, "y": 40, "w": 24, "h": 8 },
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"targets": [
{
"refId": "A",
"datasource": {
"type": "grafana-athena-datasource",
"uid": "athena"
},
"rawSQL": "SELECT wo_number, site, category, wo_status, hold_reason, mismatch FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND mismatch<>'' ORDER BY category, site",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
"database": "apm_wo_analysis",
"region": "us-east-1"
}
}
],
"options": {
"frameIndex": 0,
"showHeader": true,
"sortBy": [],
"footer": {
"show": false,
"reducer": ["sum"],
"fields": "",
"enablePagination": false
}
},
"fieldConfig": {
"defaults": {
"color": { "mode": "thresholds" },
"custom": {
"align": "left",
"cellOptions": { "type": "auto" },
"inspect": false,
"filterable": true,
"minWidth": 80,
"width": 0
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "text", "value": null }
]
}
},
"overrides": [
{
"matcher": { "id": "byName", "options": "wo_number" },
"properties": [
{ "id": "displayName", "value": "WO #" },
{ "id": "custom.width", "value": 100 }
]
},
{
"matcher": { "id": "byName", "options": "site" },
"properties": [
{ "id": "displayName", "value": "Site" },
{ "id": "custom.width", "value": 80 }
]
},
{
"matcher": { "id": "byName", "options": "category" },
"properties": [
{ "id": "displayName", "value": "Category" },
{ "id": "custom.width", "value": 160 },
{
"id": "mappings",
"value": [
{
"type": "value",
"options": {
"3rd Escalation": {
"color": "dark-red",
"index": 0
},
"2nd Escalation": {
"color": "dark-orange",
"index": 1
},
"1st Escalation": {
"color": "dark-yellow",
"index": 2
},
"SIM Ticket": {
"color": "dark-purple",
"index": 3
},
"Other Escalation": {
"color": "orange",
"index": 4
}
}
}
]
},
{
"id": "custom.cellOptions",
"value": { "type": "color-text" }
}
]
},
{
"matcher": { "id": "byName", "options": "wo_status" },
"properties": [
{ "id": "displayName", "value": "WO Status" },
{ "id": "custom.width", "value": 90 }
]
},
{
"matcher": { "id": "byName", "options": "hold_reason" },
"properties": [
{ "id": "displayName", "value": "Hold Reason" },
{ "id": "custom.width", "value": 120 }
]
},
{
"matcher": { "id": "byName", "options": "mismatch" },
"properties": [
{ "id": "displayName", "value": "Mismatch Description" },
{
"id": "custom.cellOptions",
"value": { "type": "auto", "wrapText": true }
},
{ "id": "custom.width", "value": 500 },
{ "id": "custom.minWidth", "value": 200 },
{ "id": "color", "value": { "mode": "fixed", "fixedColor": "semi-dark-orange" } }
]
}
]
}
}
]
}

View file

@ -1,12 +1,16 @@
# Athena datasource, authenticated via the EC2 instance IAM role (no static keys).
# Provisioned into /etc/grafana/provisioning/datasources/ via user-data (Phase 5).
# authType "default" = AWS SDK default credential chain, which on EC2 resolves to
# the instance role via IMDS. ("ec2_iam_role" is rejected by the plugin unless
# added to [aws] allowed_auth_providers; "default" is allowed out of the box.)
apiVersion: 1
datasources:
- name: Athena
type: grafana-athena-datasource
uid: athena
isDefault: true
jsonData:
authType: ec2_iam_role
authType: default
defaultRegion: us-east-1
catalog: AwsDataCatalog
database: apm_wo_analysis

View file

@ -30,7 +30,8 @@ import pandas as pd
import classify as clf
ANALYTICS_PREFIX = "analytics"
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.
@ -51,6 +52,7 @@ COLUMN_MATCHERS = {
}
_s3 = boto3.client("s3")
_lambda = boto3.client("lambda")
def _haiku_enabled() -> bool:
@ -177,19 +179,45 @@ def _build_summary(df: pd.DataFrame, dt: str, key: str, blank: int) -> dict:
}
def handler(event, context):
"""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"])
# 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 {"skipped": key}
return None
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()
@ -197,8 +225,8 @@ def handler(event, context):
df, blank = _build_snapshot(header, data)
if df.empty:
print("No classifiable rows (all comments blank); nothing written.")
return {"classified": 0, "blank": blank}
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
@ -213,21 +241,63 @@ def handler(event, context):
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"{ANALYTICS_PREFIX}/dt={dt}/summary.json",
Key=f"{META_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.
# 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 {
"classified": int(len(df)),
"dt": dt,
"summary_key": f"dt={dt}/summary.json",
}
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}")

View file

@ -1,3 +1,5 @@
awswrangler>=3.9.0
# awswrangler + pandas/pyarrow/numpy come from the AWS-managed SDK-for-pandas
# Lambda layer (see pipeline_stack.py) — bundling them here blows the 250 MB
# unzipped limit. The Haiku fallback uses stdlib urllib, so no anthropic SDK.
# Only openpyxl (xlsx parsing) is bundled into the function package.
openpyxl>=3.1.0
anthropic>=0.40.0

View file

@ -2,22 +2,435 @@
Two push surfaces, no App Home (see CLAUDE.md "Slack surfaces"):
- daily summary: header, vs-yesterday deltas, escalation + action/routine
fields, top sites, mismatch callout, and a 📊 Open dashboard link button to
fields, top sites, mismatch callout, and a Open dashboard link button to
Grafana. Long WO lists go in modals, never the channel (<100-block limit).
- 3rd-escalation alert: standalone, batched, one @here — suppressed entirely
on zero-3rd days (return None).
Owned by the slack-blockkit-designer agent; built in Phase 4 (docs/BUILD.md).
Slack constraints enforced here:
- Message / modal max 100 blocks.
- Actions block max 25 elements (kept to 5 per block for readability).
- section text max 3000 chars.
- Header text max 150 chars.
- No APM deep-links anywhere — Grafana only.
"""
from __future__ import annotations
# ---------------------------------------------------------------------------
# Constants
# ---------------------------------------------------------------------------
def build_daily_summary(today: dict, yesterday: dict | None) -> list[dict]:
"""Build the daily summary blocks (with vs-yesterday deltas). (Phase 4)"""
raise NotImplementedError
ESCALATION_CATEGORIES: tuple[str, ...] = (
"1st Escalation",
"2nd Escalation",
"3rd Escalation",
"SIM Ticket",
"Other Escalation",
)
# Categories that drive the "drill" buttons — non-routine and worth calling out.
# Order determines button priority when we cap at 5 per actions block.
DRILL_CATEGORY_PRIORITY: tuple[str, ...] = (
"3rd Escalation",
"2nd Escalation",
"1st Escalation",
"SIM Ticket",
"Vendor No-Show",
"Awaiting Scheduling",
"Report / Docs Needed",
"Awaiting Report / Invoice",
"Awaiting Vendor / Parts",
"Status Inquiry",
"Other Escalation",
)
# Max WOs to render inline before the "+M more" overflow note in the alert.
ALERT_MAX_WOS = 20
# Max blocks reserved for WO rows in the modal (leaves headroom for
# title/footer blocks). Each WO gets 1 section block.
MODAL_HEADER_BLOCKS = 1 # title is not a block; we add 1 context block as preamble
MODAL_FOOTER_BLOCKS = 1 # overflow context block (may not be emitted)
MODAL_MAX_WO_BLOCKS = 100 - MODAL_HEADER_BLOCKS - MODAL_FOOTER_BLOCKS # = 98
def build_escalation_alert(today: dict) -> list[dict] | None:
"""Build the batched 3rd-escalation alert, or None when count is 0. (Phase 4)"""
raise NotImplementedError
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _divider() -> dict:
return {"type": "divider"}
def _header(text: str) -> dict:
# Header block text is plain_text, max 150 chars.
return {
"type": "header",
"text": {"type": "plain_text", "text": text[:150], "emoji": True},
}
def _section(text: str) -> dict:
# Section text is mrkdwn, max 3000 chars.
return {
"type": "section",
"text": {"type": "mrkdwn", "text": text[:3000]},
}
def _context(text: str) -> dict:
return {
"type": "context",
"elements": [{"type": "mrkdwn", "text": text[:3000]}],
}
def _fields_section(fields: list[str]) -> dict:
"""Section with up to 10 fields (Slack limit), each <=2000 chars."""
return {
"type": "section",
"fields": [{"type": "mrkdwn", "text": f[:2000]} for f in fields[:10]],
}
def _button(text: str, action_id: str, value: str) -> dict:
"""Interactive button (opens modal via interactivity handler)."""
return {
"type": "button",
"text": {"type": "plain_text", "text": text, "emoji": False},
"action_id": action_id,
"value": value,
}
def _link_button(text: str, url: str) -> dict:
"""Link button — opens URL directly, no action_id."""
return {
"type": "button",
"text": {"type": "plain_text", "text": text, "emoji": True},
"url": url,
}
def _actions(elements: list[dict]) -> dict:
"""Actions block; Slack limit 25 elements, we keep <=5 per block."""
return {"type": "actions", "elements": elements[:25]}
def _delta_str(now: int, prev: int) -> str:
"""Format a vs-yesterday delta as ▲+N or ▼-N or ~0."""
diff = now - prev
if diff > 0:
return f"▲+{diff}"
if diff < 0:
return f"▼{diff}"
return "~0"
def _fmt_date(dt: str) -> str:
"""Format 'YYYY-MM-DD' as 'May 28, 2026' for display."""
try:
from datetime import date
d = date.fromisoformat(dt)
# %-d is Linux-specific; strip the leading zero portably instead.
return f"{d.strftime('%b')} {d.day}, {d.year}"
except Exception:
return dt
# ---------------------------------------------------------------------------
# Public surface builders
# ---------------------------------------------------------------------------
def build_daily_summary(
today: dict,
yesterday: dict | None,
dashboard_url: str,
) -> list[dict]:
"""Build the daily summary post blocks.
Returns a list of Block Kit block dicts ready to pass to ``blocks=`` in a
``chat.postMessage`` call. Always stays well under 100 blocks.
Args:
today: The contents of today's ``summary.json``.
yesterday: The contents of yesterday's ``summary.json``, or None when
there is no previous snapshot (first run, etc.).
dashboard_url: The Grafana dashboard URL for the link button.
"""
blocks: list[dict] = []
dt = today.get("dt", "")
date_label = _fmt_date(dt) if dt else "Today"
# ------------------------------------------------------------------
# 1. Header
# ------------------------------------------------------------------
blocks.append(_header(f"APM Work Orders — {date_label}"))
# ------------------------------------------------------------------
# 2. Context line: classified total + vs-yesterday deltas
# ------------------------------------------------------------------
classified = today.get("classified_total", 0)
blank = today.get("blank_comment_rows", 0)
escalation = today.get("escalation_total", 0)
action = today.get("action_needed", 0)
routine = today.get("routine", 0)
if yesterday is not None:
delta_classified = _delta_str(classified, yesterday.get("classified_total", 0))
delta_escalation = _delta_str(escalation, yesterday.get("escalation_total", 0))
delta_action = _delta_str(action, yesterday.get("action_needed", 0))
context_text = (
f"*{classified}* WOs classified ({blank} blank-comment rows excluded) "
f"| Escalations {delta_escalation} | Action-needed {delta_action} "
f"| vs yesterday: total {delta_classified}"
)
else:
context_text = (
f"*{classified}* WOs classified ({blank} blank-comment rows excluded) "
f"| No prior day for delta comparison"
)
blocks.append(_section(context_text))
blocks.append(_divider())
# ------------------------------------------------------------------
# 3. Escalation breakdown
# ------------------------------------------------------------------
category_counts: dict[str, int] = today.get("category_counts", {})
third_count = today.get(
"third_escalation_count", category_counts.get("3rd Escalation", 0)
)
esc_lines: list[str] = []
for cat in ESCALATION_CATEGORIES:
count = category_counts.get(cat, 0)
if count == 0:
continue
prefix = "*" if cat == "3rd Escalation" else ""
suffix = "*" if cat == "3rd Escalation" else ""
esc_lines.append(f"{prefix}{cat}: {count}{suffix}")
if esc_lines:
blocks.append(_section("*Escalations*\n" + "\n".join(esc_lines)))
else:
blocks.append(_section("*Escalations*\nNone today"))
# ------------------------------------------------------------------
# 4. Action-needed vs routine fields
# ------------------------------------------------------------------
blocks.append(
_fields_section(
[
f"*Action-needed*\n{action}",
f"*Routine*\n{routine}",
f"*3rd Escalations*\n{third_count}",
f"*Total classified*\n{classified}",
]
)
)
blocks.append(_divider())
# ------------------------------------------------------------------
# 5. Top 5 sites
# ------------------------------------------------------------------
top_sites: list[dict] = today.get("top_sites", [])[:5]
if top_sites:
site_lines = "\n".join(
f"{i + 1}. *{s['site']}* — {s['count']} WOs"
for i, s in enumerate(top_sites)
)
blocks.append(_section(f"*Top sites*\n{site_lines}"))
blocks.append(_divider())
# ------------------------------------------------------------------
# 6. Mismatch callout (only when mismatches present)
# ------------------------------------------------------------------
mismatches: list[dict] = today.get("mismatches", [])
if mismatches:
mm_count = len(mismatches)
blocks.append(
_section(
f":warning: *{mm_count} mismatch{'es' if mm_count != 1 else ''} "
f"flagged* — comment intent contradicts structured state. "
f"See the Grafana mismatch panel for details."
)
)
blocks.append(_divider())
# ------------------------------------------------------------------
# 7. Drill buttons (action_id="drill_category") + dashboard link button
# ------------------------------------------------------------------
# Pick the most-common non-routine categories up to 4 slots; the 5th slot
# is always the dashboard link button, keeping each actions block <=5.
drill_cats: list[tuple[str, int]] = []
for cat in DRILL_CATEGORY_PRIORITY:
count = category_counts.get(cat, 0)
if count > 0:
drill_cats.append((cat, count))
if len(drill_cats) >= 4:
break
button_elements: list[dict] = []
for cat, count in drill_cats:
short_label = cat.replace(" Escalation", " Esc.").replace("Awaiting ", "")
short_label = short_label[:20] # button text kept concise
# action_id must be UNIQUE per message (Slack rejects duplicates), so
# qualify it with the category; the interactions handler matches on the
# "drill_category:" prefix and reads the filter value from `value`.
button_elements.append(
_button(f"{short_label} ({count})", f"drill_category:{cat}", cat)
)
# Dashboard link button always present (no action_id — url button).
button_elements.append(_link_button("Open dashboard", dashboard_url))
blocks.append(_actions(button_elements))
# ------------------------------------------------------------------
# 8. Footer context
# ------------------------------------------------------------------
generated_at = today.get("generated_at", "")
footer_parts = [f"Generated {generated_at}" if generated_at else ""]
footer_parts.append(f"dt={dt}")
blocks.append(_context(" | ".join(p for p in footer_parts if p)))
# Safety check — should never fire in practice given the bounded sections
# above, but surface a warning in the footer rather than silently truncate.
if len(blocks) > 95:
# Trim from the middle, preserve header and footer.
over = len(blocks) - 95
blocks = blocks[:3] + blocks[3 + over :]
blocks.append(
_context(f"+{over} block(s) trimmed — view full breakdown in Grafana.")
)
return blocks
def build_escalation_alert(escalations: list[dict]) -> list[dict] | None:
"""Build the batched 3rd-escalation alert blocks.
Zero-3rd suppression: returns None when the list is empty. The caller must
check for None before posting — do NOT post a message with no content.
Args:
escalations: Rows from ``details.json`` pre-filtered to
``category == "3rd Escalation"``.
Returns:
A list of Block Kit block dicts, or None.
"""
if not escalations:
return None
blocks: list[dict] = []
# ------------------------------------------------------------------
# 1. Header
# ------------------------------------------------------------------
count = len(escalations)
blocks.append(
_header(f"3rd Escalation Alert — {count} WO{'s' if count != 1 else ''}")
)
# ------------------------------------------------------------------
# 2. ONE @here mention (never one ping per WO)
# ------------------------------------------------------------------
blocks.append(
_section(
f"<!here> *{count} work order{'s' if count != 1 else ''} "
f"at 3rd escalation* require immediate attention."
)
)
blocks.append(_divider())
# ------------------------------------------------------------------
# 3. Compact WO list, capped at ALERT_MAX_WOS
# ------------------------------------------------------------------
visible = escalations[:ALERT_MAX_WOS]
overflow = len(escalations) - len(visible)
for wo in visible:
wo_num = wo.get("wo_number", "???")
site = wo.get("site", "—")
desc = wo.get("wo_description", "")
# Truncate description to keep the block tidy.
if len(desc) > 80:
desc = desc[:77] + "..."
blocks.append(_section(f"*{wo_num}* — {site} — {desc}"))
if overflow > 0:
blocks.append(
_context(f"+{overflow} more — see the Grafana dashboard for the full list.")
)
return blocks
def build_wo_modal(title: str, wos: list[dict], dashboard_url: str) -> dict:
"""Build a ``views.open`` payload (modal) listing work orders.
Caps to Slack's 100-block modal limit by rendering the first chunk of WOs
and appending a "+M more" overflow context block when the list is too long.
No APM deep-links — Grafana only.
Args:
title: Modal title (max 24 chars enforced by Slack; truncated here).
wos: List of per-WO dicts from ``details.json``.
dashboard_url: The Grafana dashboard URL for the overflow note link.
Returns:
A dict suitable for passing to ``views.open`` as the ``view=`` arg.
"""
# Slack enforces a 24-char modal title limit.
title_text = title[:24]
modal_blocks: list[dict] = []
# Each WO is one section block. Reserve space for the optional overflow
# context block.
max_wo = MODAL_MAX_WO_BLOCKS - 1 # one slot reserved for potential overflow
visible = wos[:max_wo]
overflow_count = len(wos) - len(visible)
for wo in visible:
wo_num = wo.get("wo_number", "???")
site = wo.get("site", "—")
dept = wo.get("department", "—")
category = wo.get("category", "—")
snippet = wo.get("last_comment", "")
if len(snippet) > 150:
snippet = snippet[:147] + "..."
text_parts = [
f"*{wo_num}*",
f"_{site}_ | {dept} | {category}",
]
if snippet:
text_parts.append(snippet)
modal_blocks.append(_section("\n".join(text_parts)))
if overflow_count > 0:
modal_blocks.append(
_context(
f"+{overflow_count} more — <{dashboard_url}|view the full list in Grafana>."
)
)
elif not modal_blocks:
modal_blocks.append(_section("No work orders to display."))
return {
"type": "modal",
"title": {"type": "plain_text", "text": title_text, "emoji": False},
"close": {"type": "plain_text", "text": "Close", "emoji": False},
"blocks": modal_blocks,
}

View file

@ -1,16 +1,59 @@
"""Slack post + alert Lambda: daily summary and batched 3rd-escalation alert.
"""Slack post Lambda: daily summary + batched 3rd-escalation alert.
Invoked after the classifier finishes. Reads ``analytics/dt=today/summary.json``
and the yesterday partition (for deltas), posts the daily summary to the WO
channel via the reused bot token (SSM, same pattern as payments-dashboard), and
— only if the 3rd-escalation count > 0 — posts the standalone batched alert.
Async-invoked by the classifier with ``{"dt": "YYYY-MM-DD"}`` once the day's
snapshot is written. Reads today's (and yesterday's, for deltas) ``summary.json``
from ``analytics/``, posts the daily summary to the WO channel via the reused bot
token (Secrets Manager), and — only when the 3rd-escalation count > 0 — posts the
standalone batched alert (reading ``details.json`` for the 3rd-escalation rows).
Implemented in Phase 4 (docs/BUILD.md).
Runtime: Python 3.12, ARM64. No App Home — two push surfaces only (see CLAUDE.md).
"""
from __future__ import annotations
from datetime import datetime, timedelta, timezone
import blockkit
import slackio
def _yesterday(dt: str) -> str:
return (datetime.strptime(dt, "%Y-%m-%d") - timedelta(days=1)).strftime("%Y-%m-%d")
def handler(event, context):
"""Lambda entry point. (Phase 4)"""
raise NotImplementedError
"""Post the daily summary and (conditionally) the 3rd-escalation alert."""
dt = (event or {}).get("dt") or datetime.now(timezone.utc).strftime("%Y-%m-%d")
today = slackio.read_meta_json(dt, "summary.json")
if today is None:
print(f"No summary.json for dt={dt}; nothing to post.")
return {"posted": False, "reason": "no summary", "dt": dt}
yesterday = slackio.read_meta_json(_yesterday(dt), "summary.json")
client = slackio.web_client()
channel = slackio.channel_id()
dashboard_url = slackio.get_dashboard_url()
summary_blocks = blockkit.build_daily_summary(today, yesterday, dashboard_url)
client.chat_postMessage(
channel=channel, blocks=summary_blocks, text=f"APM Work Orders — {dt}"
)
posted = {"summary": True, "alert": False}
# Standalone batched alert — only when there are 3rd escalations. Pull the
# rows from details.json so the alert can name the WOs.
if today.get("third_escalation_count", 0) > 0:
details = slackio.read_meta_json(dt, "details.json") or []
thirds = [d for d in details if d.get("category") == "3rd Escalation"]
alert_blocks = blockkit.build_escalation_alert(thirds)
if alert_blocks:
client.chat_postMessage(
channel=channel,
blocks=alert_blocks,
text=f"🚨 {len(thirds)} 3rd-escalation work orders",
)
posted["alert"] = True
print(f"Posted dt={dt}: {posted}")
return {"posted": True, "dt": dt, **posted}

View file

@ -0,0 +1,72 @@
"""Slack interactivity Lambda: drill-down modals behind API Gateway.
Fronted by the ``apm-wo.seahaven.com`` HTTP API (Slack interactivity request URL).
Verifies the Slack request signature, then for a category/site drill button reads
the day's ``details.json``, filters it, and opens a WO-list modal via ``views.open``
within Slack's 3-second ``trigger_id`` window. WO rows link to Grafana only — no
APM deep-links (see Phase 4 decisions).
"""
from __future__ import annotations
import base64
import json
from datetime import datetime, timezone
from urllib.parse import parse_qs
import blockkit
import slackio
_FILTER_FIELD = {"drill_category": "category", "drill_site": "site"}
def _raw_body(event: dict) -> str:
body = event.get("body") or ""
if event.get("isBase64Encoded"):
body = base64.b64decode(body).decode("utf-8")
return body
def _headers(event: dict) -> dict:
# API Gateway HTTP API lowercases header names; be defensive either way.
return {k.lower(): v for k, v in (event.get("headers") or {}).items()}
def handler(event, context):
"""Verify the Slack signature and open the requested drill-down modal."""
body = _raw_body(event)
headers = _headers(event)
timestamp = headers.get("x-slack-request-timestamp", "")
signature = headers.get("x-slack-signature", "")
if not slackio.verify_signature(body, timestamp, signature):
print("Rejected: invalid Slack signature.")
return {"statusCode": 401, "body": "invalid signature"}
parsed = parse_qs(body)
payload = json.loads(parsed.get("payload", ["{}"])[0])
if payload.get("type") != "block_actions":
return {"statusCode": 200, "body": ""}
action = (payload.get("actions") or [{}])[0]
# action_id is qualified for Slack uniqueness, e.g. "drill_category:Report /
# Docs Needed" — match on the prefix before ":"; the filter value is in `value`.
action_kind = (action.get("action_id") or "").split(":", 1)[0]
field = _FILTER_FIELD.get(action_kind)
value = action.get("value")
trigger_id = payload.get("trigger_id")
if not field or not value or not trigger_id:
return {"statusCode": 200, "body": ""}
# The daily post is same-day; default the drill to today's snapshot.
dt = datetime.now(timezone.utc).strftime("%Y-%m-%d")
details = slackio.read_meta_json(dt, "details.json") or []
wos = [d for d in details if d.get(field) == value]
view = blockkit.build_wo_modal(
title=f"{value} ({len(wos)})",
wos=wos,
dashboard_url=slackio.get_dashboard_url(),
)
slackio.web_client().views_open(trigger_id=trigger_id, view=view)
return {"statusCode": 200, "body": ""}

View file

@ -0,0 +1,71 @@
"""Shared Slack + AWS I/O for the post and interactions Lambdas.
Keeps the Block Kit builders (blockkit.py) pure: everything that touches the
network or AWS lives here. Slack credentials are a single Secrets Manager secret
``{ botToken, signingSecret, channelId }``; the Grafana dashboard URL is
operational config in SSM (editable without a redeploy). Both are cached for the
life of the execution environment.
"""
from __future__ import annotations
import json
import os
import boto3
from slack_sdk import WebClient
from slack_sdk.signature import SignatureVerifier
_secrets = boto3.client("secretsmanager")
_ssm = boto3.client("ssm")
_s3 = boto3.client("s3")
_SECRET_NAME = os.environ["SLACK_SECRET_NAME"]
_DASHBOARD_PARAM = os.environ["DASHBOARD_URL_PARAM"]
_BUCKET = os.environ["ANALYTICS_BUCKET"]
# Lazily-populated caches (warm across invocations in the same container).
_creds: dict | None = None
_dashboard_url: str | None = None
def get_credentials() -> dict:
"""Return the Slack creds dict: ``botToken``, ``signingSecret``, ``channelId``."""
global _creds
if _creds is None:
raw = _secrets.get_secret_value(SecretId=_SECRET_NAME)["SecretString"]
_creds = json.loads(raw)
return _creds
def get_dashboard_url() -> str:
"""Grafana dashboard URL for the 📊 button / modal overflow links (SSM)."""
global _dashboard_url
if _dashboard_url is None:
_dashboard_url = _ssm.get_parameter(Name=_DASHBOARD_PARAM)["Parameter"]["Value"]
return _dashboard_url
def web_client() -> WebClient:
return WebClient(token=get_credentials()["botToken"])
def channel_id() -> str:
return get_credentials()["channelId"]
def verify_signature(body: str, timestamp: str, signature: str) -> bool:
"""Validate a Slack request signature (HMAC + 5-minute replay window)."""
verifier = SignatureVerifier(signing_secret=get_credentials()["signingSecret"])
return verifier.is_valid(body=body, timestamp=timestamp, signature=signature)
def read_meta_json(dt: str, name: str):
"""Read ``meta/dt=<dt>/<name>`` as JSON, or None if absent. (summary.json /
details.json live under meta/, kept out of the Athena table's analytics/ prefix.)"""
key = f"meta/dt={dt}/{name}"
try:
obj = _s3.get_object(Bucket=_BUCKET, Key=key)
except _s3.exceptions.NoSuchKey:
return None
return json.loads(obj["Body"].read())

41
slack/manifest.yaml Normal file
View file

@ -0,0 +1,41 @@
# Slack app manifest for "APM Work Orders".
#
# IMPORTANT: The request_url host (apm-wo.seahaven.com) must match the deployed
# API Gateway custom domain. If the domain changes, update interactivity.request_url
# here and redeploy the app via the Slack API or the App Configuration page.
#
# The bot user must be invited to the WO channel before it can post:
# /invite @APM Work Orders
#
# To apply this manifest: Slack App Configuration -> "App Manifest" tab -> paste.
_metadata:
major_version: 1
minor_version: 1
display_information:
name: APM Work Orders
description: Daily APM work-order analysis — escalation alerts and summary posts for Sea Haven facility ops.
background_color: "#1a1a2e"
features:
bot_user:
display_name: APM Work Orders
always_online: true
oauth_config:
scopes:
bot:
- chat:write
# chat:write.public allows posting to public channels without an explicit
# /invite. Remove if the WO channel is private and an invite is preferred.
- chat:write.public
settings:
interactivity:
is_enabled: true
# Handler for drill-category buttons and modal opens.
request_url: https://apm-wo.seahaven.com/slack/interactions
org_deploy_enabled: false
socket_mode_enabled: false
token_rotation_enabled: false

454
tests/test_blockkit.py Normal file
View file

@ -0,0 +1,454 @@
"""Offline unit tests for the Block Kit surface builders.
No AWS, no network, no slack_sdk. All fixtures are plain dicts.
Run with the repo venv:
./.venv/bin/python -m pytest tests/test_blockkit.py -q
"""
from __future__ import annotations
import sys
from pathlib import Path
sys.path.insert(
0, str(Path(__file__).resolve().parent.parent / "lambdas" / "slack_post")
)
import blockkit # noqa: E402
from blockkit import ( # noqa: E402
build_daily_summary,
build_escalation_alert,
build_wo_modal,
)
# ---------------------------------------------------------------------------
# Shared fixtures
# ---------------------------------------------------------------------------
DASHBOARD_URL = "https://grafana.seahaven.internal/d/apm-wo"
CATEGORY_COUNTS_TYPICAL: dict[str, int] = {
"3rd Escalation": 4,
"2nd Escalation": 8,
"1st Escalation": 12,
"SIM Ticket": 2,
"Vendor No-Show": 3,
"Awaiting Scheduling": 25,
"Report / Docs Needed": 18,
"Awaiting Report / Invoice": 7,
"Awaiting Vendor / Parts": 5,
"Status Inquiry": 6,
"Other Escalation": 1,
"Schedule Confirmed": 90,
"Weekly WO Scheduled": 40,
"Completed / Pending Close": 22,
"Acknowledgement / No-op": 10,
"Cancelled": 5,
"On Hold": 3,
"Rescheduled": 4,
"Avetta Project Created": 2,
"Other": 15,
}
TODAY_SUMMARY: dict = {
"dt": "2026-05-28",
"classified_total": 282,
"blank_comment_rows": 68,
"category_counts": CATEGORY_COUNTS_TYPICAL,
"escalation_total": 27,
"third_escalation_count": 4,
"action_needed": 91,
"routine": 191,
"top_sites": [
{"site": "ABQ5", "count": 45},
{"site": "ACY9", "count": 38},
{"site": "BOS1", "count": 31},
{"site": "DFW7", "count": 28},
{"site": "LAX9", "count": 22},
],
"mismatches": [
{
"wo_number": "WO-001",
"category": "Completed / Pending Close",
"mismatch": "comment claims completion while on REPORT hold",
},
{
"wo_number": "WO-002",
"category": "Schedule Confirmed",
"mismatch": "comment confirms schedule while on SCHEDULING hold",
},
],
"generated_at": "2026-05-28T06:00:00Z",
}
YESTERDAY_SUMMARY: dict = {
"dt": "2026-05-27",
"classified_total": 270,
"blank_comment_rows": 80,
"category_counts": {},
"escalation_total": 30,
"third_escalation_count": 6,
"action_needed": 95,
"routine": 175,
"top_sites": [],
"mismatches": [],
"generated_at": "2026-05-27T06:00:00Z",
}
def _make_wo(i: int) -> dict:
return {
"wo_number": f"WO-{i:04d}",
"wo_description": f"Repair HVAC unit at dock {i}",
"site": "ABQ5",
"department": "SSP",
"category": "3rd Escalation",
"last_comment": f"3rd attempt to schedule vendor for dock {i} repair.",
"is_escalation": True,
"is_action": True,
"mismatch": None,
}
# ---------------------------------------------------------------------------
# build_daily_summary
# ---------------------------------------------------------------------------
class TestDailySummaryDeltas:
def test_with_yesterday_shows_deltas(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
# Collect all text from section/context blocks.
all_text = _extract_text(blocks)
# classified went 270 -> 282 (+12), escalation 30 -> 27 (-3).
assert "▲+12" in all_text, "expected classified delta ▲+12"
assert "▼-3" in all_text, "expected escalation delta ▼-3"
def test_without_yesterday_omits_deltas(self):
"""yesterday=None must not raise and must omit delta symbols."""
blocks = build_daily_summary(TODAY_SUMMARY, None, DASHBOARD_URL)
all_text = _extract_text(blocks)
assert "▲" not in all_text
assert "▼" not in all_text
assert "No prior day" in all_text
def test_without_yesterday_returns_blocks(self):
blocks = build_daily_summary(TODAY_SUMMARY, None, DASHBOARD_URL)
assert isinstance(blocks, list)
assert len(blocks) > 0
def test_zero_delta_shows_tilde(self):
yesterday_same = {**YESTERDAY_SUMMARY, "escalation_total": 27}
blocks = build_daily_summary(TODAY_SUMMARY, yesterday_same, DASHBOARD_URL)
all_text = _extract_text(blocks)
assert "~0" in all_text
class TestDailySummaryStructure:
def test_has_header_block(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
assert any(b.get("type") == "header" for b in blocks)
def test_action_ids_are_unique(self):
# Slack rejects a message with duplicate action_ids across its elements
# (regression: all drill buttons once shared action_id "drill_category").
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
action_ids = [
el["action_id"]
for b in blocks
for el in b.get("elements", [])
if isinstance(el, dict) and "action_id" in el
]
assert len(action_ids) == len(set(action_ids)), (
f"duplicate action_id(s): {action_ids}"
)
def test_header_contains_date(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
header = next(b for b in blocks if b.get("type") == "header")
assert "2026" in header["text"]["text"] or "May" in header["text"]["text"]
def test_mismatch_callout_present_when_mismatches_nonzero(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
all_text = _extract_text(blocks)
assert "mismatch" in all_text.lower()
def test_mismatch_callout_absent_when_no_mismatches(self):
today_no_mm = {**TODAY_SUMMARY, "mismatches": []}
blocks = build_daily_summary(today_no_mm, YESTERDAY_SUMMARY, DASHBOARD_URL)
all_text = _extract_text(blocks)
assert "mismatch" not in all_text.lower()
def test_dashboard_link_button_url(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
url = _find_link_button_url(blocks)
assert url == DASHBOARD_URL, f"expected {DASHBOARD_URL!r}, got {url!r}"
def test_dashboard_button_has_no_action_id(self):
"""The Open-dashboard button must be a url button, not an action button."""
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
for block in blocks:
if block.get("type") != "actions":
continue
for elem in block.get("elements", []):
if elem.get("url") == DASHBOARD_URL:
assert "action_id" not in elem, (
"dashboard link button must not carry an action_id"
)
def test_drill_buttons_have_action_id(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
drill_found = False
for block in blocks:
if block.get("type") != "actions":
continue
for elem in block.get("elements", []):
# action_id is qualified for uniqueness: "drill_category:<cat>".
if str(elem.get("action_id", "")).startswith("drill_category:"):
drill_found = True
assert "value" in elem
assert drill_found, "expected at least one drill_category action button"
def test_3rd_escalation_highlighted(self):
blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL)
all_text = _extract_text(blocks)
# 3rd Escalation should appear in bold mrkdwn.
assert "*3rd Escalation" in all_text or "3rd Escalation*" in all_text
class TestDailySummaryBlockLimit:
def _large_summary(self) -> dict:
"""Summary with 25+ categories, 30 sites, and 200 mismatches."""
big_cats: dict[str, int] = {f"Category {i}": i + 1 for i in range(30)}
big_cats.update(CATEGORY_COUNTS_TYPICAL)
big_sites = [{"site": f"S{i:03d}", "count": 30 - i} for i in range(30)]
big_mismatches = [
{"wo_number": f"WO-{i:04d}", "category": "X", "mismatch": "test mm"}
for i in range(200)
]
return {
**TODAY_SUMMARY,
"category_counts": big_cats,
"top_sites": big_sites,
"mismatches": big_mismatches,
}
def test_stays_under_100_blocks_large_input(self):
big = self._large_summary()
blocks = build_daily_summary(big, YESTERDAY_SUMMARY, DASHBOARD_URL)
assert len(blocks) < 100, f"daily summary exceeded 100 blocks: {len(blocks)}"
def test_reports_block_count(self, capsys):
big = self._large_summary()
blocks = build_daily_summary(big, YESTERDAY_SUMMARY, DASHBOARD_URL)
with capsys.disabled():
print(f"\n[test] daily summary (large input): {len(blocks)} blocks")
# ---------------------------------------------------------------------------
# build_escalation_alert
# ---------------------------------------------------------------------------
class TestEscalationAlert:
def test_returns_none_on_empty_list(self):
assert build_escalation_alert([]) is None
def test_returns_blocks_on_non_empty(self):
escalations = [_make_wo(i) for i in range(3)]
result = build_escalation_alert(escalations)
assert result is not None
assert isinstance(result, list)
assert len(result) > 0
def test_exactly_one_here_mention(self):
escalations = [_make_wo(i) for i in range(5)]
blocks = build_escalation_alert(escalations)
assert blocks is not None
all_text = _extract_text(blocks)
here_count = all_text.count("<!here>")
assert here_count == 1, f"expected exactly 1 <!here>, found {here_count}"
def test_here_in_section_not_header(self):
"""@here must be in a section block, not the header."""
escalations = [_make_wo(i) for i in range(2)]
blocks = build_escalation_alert(escalations)
assert blocks is not None
for block in blocks:
if block.get("type") == "header":
header_text = block["text"]["text"]
assert "<!here>" not in header_text
def test_no_here_on_empty(self):
result = build_escalation_alert([])
assert result is None
def test_overflow_note_when_many_wos(self):
"""When list exceeds ALERT_MAX_WOS, a '+M more' context block is added."""
many = [_make_wo(i) for i in range(blockkit.ALERT_MAX_WOS + 10)]
blocks = build_escalation_alert(many)
assert blocks is not None
all_text = _extract_text(blocks)
assert "more" in all_text.lower()
def test_no_overflow_note_when_within_limit(self):
few = [_make_wo(i) for i in range(blockkit.ALERT_MAX_WOS - 1)]
blocks = build_escalation_alert(few)
assert blocks is not None
all_text = _extract_text(blocks)
# Should not have a "+N more" line.
assert " more —" not in all_text
def test_single_wo_grammatically_correct(self):
"""'WO' (singular) when count is 1."""
blocks = build_escalation_alert([_make_wo(1)])
assert blocks is not None
header = next(b for b in blocks if b.get("type") == "header")
assert "WOs" not in header["text"]["text"], "expected singular 'WO' not 'WOs'"
def test_stays_under_100_blocks(self):
huge = [_make_wo(i) for i in range(500)]
blocks = build_escalation_alert(huge)
assert blocks is not None
assert len(blocks) < 100, f"alert exceeded 100 blocks: {len(blocks)}"
def test_wo_number_appears_in_blocks(self):
escalations = [_make_wo(42)]
blocks = build_escalation_alert(escalations)
assert blocks is not None
all_text = _extract_text(blocks)
assert "WO-0042" in all_text
# ---------------------------------------------------------------------------
# build_wo_modal
# ---------------------------------------------------------------------------
class TestWoModal:
def test_returns_dict_with_modal_type(self):
modal = build_wo_modal("Test Modal", [_make_wo(1)], DASHBOARD_URL)
assert isinstance(modal, dict)
assert modal.get("type") == "modal"
def test_modal_has_close_button(self):
modal = build_wo_modal("Test", [_make_wo(1)], DASHBOARD_URL)
assert "close" in modal
assert modal["close"]["type"] == "plain_text"
def test_title_truncated_to_24_chars(self):
long_title = "A" * 50
modal = build_wo_modal(long_title, [_make_wo(1)], DASHBOARD_URL)
assert len(modal["title"]["text"]) <= 24
def test_500_wos_truncated_to_under_100_blocks(self, capsys):
wos = [_make_wo(i) for i in range(500)]
modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL)
block_count = len(modal["blocks"])
with capsys.disabled():
print(f"\n[test] 500-WO modal: {block_count} blocks")
assert block_count <= 100, (
f"modal with 500 WOs has {block_count} blocks, exceeds 100"
)
def test_500_wos_has_overflow_note(self):
wos = [_make_wo(i) for i in range(500)]
modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
assert "more" in all_text.lower(), "expected '+M more' overflow note"
def test_overflow_note_contains_dashboard_url(self):
wos = [_make_wo(i) for i in range(500)]
modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
assert DASHBOARD_URL in all_text, "overflow note must link to the dashboard URL"
def test_overflow_count_correct(self):
"""The overflow note's M value must equal len(wos) - rendered_count."""
n = 500
wos = [_make_wo(i) for i in range(n)]
modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL)
blocks = modal["blocks"]
# Last block should be context with the overflow note.
last_text = _extract_text([blocks[-1]])
# The overflow note pattern is "+M more".
import re
match = re.search(r"\+(\d+) more", last_text)
assert match, f"no '+M more' pattern in last block: {last_text!r}"
reported_overflow = int(match.group(1))
# Block count: total blocks = visible WO blocks + 1 overflow block.
rendered_wo_blocks = len(blocks) - 1
assert reported_overflow == n - rendered_wo_blocks, (
f"overflow note says +{reported_overflow} but "
f"{n} - {rendered_wo_blocks} = {n - rendered_wo_blocks}"
)
def test_empty_wo_list_returns_placeholder(self):
modal = build_wo_modal("Empty", [], DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
assert "no work orders" in all_text.lower()
def test_wo_number_appears_in_blocks(self):
modal = build_wo_modal("Test", [_make_wo(7)], DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
assert "WO-0007" in all_text
def test_snippet_truncated_to_150_chars(self):
wo = _make_wo(1)
wo["last_comment"] = "X" * 300
modal = build_wo_modal("Test", [wo], DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
# The snippet in the text should not be 300 X's.
assert "X" * 200 not in all_text
def test_no_apm_deep_links(self):
"""Modal must not contain any APM application deep-links."""
wos = [_make_wo(i) for i in range(10)]
modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL)
all_text = _extract_text(modal["blocks"])
# APM URLs would start with typical patterns — confirm none present.
assert "apm://" not in all_text
assert "app.apm" not in all_text
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _extract_text(blocks: list[dict]) -> str:
"""Collect all text strings from blocks for assertion convenience."""
parts: list[str] = []
for block in blocks:
btype = block.get("type")
if btype == "header":
parts.append(block.get("text", {}).get("text", ""))
elif btype == "section":
text_obj = block.get("text")
if text_obj:
parts.append(text_obj.get("text", ""))
for field in block.get("fields", []):
parts.append(field.get("text", ""))
elif btype == "context":
for elem in block.get("elements", []):
parts.append(elem.get("text", ""))
elif btype == "actions":
for elem in block.get("elements", []):
txt = elem.get("text", {})
if isinstance(txt, dict):
parts.append(txt.get("text", ""))
return "\n".join(parts)
def _find_link_button_url(blocks: list[dict]) -> str | None:
"""Return the URL from the first link button found in any actions block."""
for block in blocks:
if block.get("type") != "actions":
continue
for elem in block.get("elements", []):
if "url" in elem:
return elem["url"]
return None

181
tests/test_grafana_synth.py Normal file
View file

@ -0,0 +1,181 @@
"""Synth-level assertions for the Grafana stack (Phase 5).
Synthesizes ``apm-wo-analysis-grafana`` and asserts the security posture that
can't be eyeballed: the ALB only admits the office CIDRs on 443 (never
0.0.0.0/0), the instance only takes traffic from the ALB SG, the instance role
carries no static keys and only scoped Athena/Glue-read/S3 access, the root
volume is gp3 + retained, a daily DLM backup exists, and grafana.seahaven.com
aliases the ALB. No AWS, no Docker (bundling skipped).
Run with the repo venv:
python -m pytest tests/test_grafana_synth.py -q
"""
import json
import sys
from pathlib import Path
import aws_cdk as cdk
from aws_cdk.assertions import Match, Template
CDK_DIR = Path(__file__).resolve().parents[1] / "cdk"
sys.path.insert(0, str(CDK_DIR))
from stacks.grafana_stack import GrafanaStack # noqa: E402
OFFICE_CIDRS = {"47.21.61.4/32", "96.250.164.146/32"}
def _cdk_context() -> dict:
ctx = json.loads((CDK_DIR / "cdk.json").read_text())["context"]
ctx["aws:cdk:bundling-stacks"] = []
return ctx
def _template() -> Template:
app = cdk.App(context=_cdk_context())
stack = GrafanaStack(
app,
"apm-wo-analysis-grafana",
env=cdk.Environment(account="328440206208", region="us-east-1"),
)
return Template.from_stack(stack)
def _all_cidr_ingress(t: Template):
"""Every CIDR-based ingress rule, inline on SGs and standalone, as
(cidr, from_port, to_port) tuples."""
rules = []
for sg in t.find_resources("AWS::EC2::SecurityGroup").values():
for r in sg["Properties"].get("SecurityGroupIngress", []):
if "CidrIp" in r:
rules.append((r["CidrIp"], r.get("FromPort"), r.get("ToPort")))
for ing in t.find_resources("AWS::EC2::SecurityGroupIngress").values():
p = ing["Properties"]
if "CidrIp" in p:
rules.append((p["CidrIp"], p.get("FromPort"), p.get("ToPort")))
return rules
def test_alb_only_admits_office_cidrs_on_443():
rules = _all_cidr_ingress(_template())
cidrs_443 = {c for c, fp, tp in rules if fp == 443 and tp == 443}
assert cidrs_443 == OFFICE_CIDRS, (
f"443 ingress should be office-only, got {cidrs_443}"
)
# Nothing anywhere may be open to the world.
assert all(c != "0.0.0.0/0" for c, _, _ in rules), "found a 0.0.0.0/0 ingress"
def test_instance_only_reachable_from_alb_on_3000():
# The instance SG ingress on 3000 is a SourceSecurityGroup rule, not a CIDR.
_template().has_resource_properties(
"AWS::EC2::SecurityGroupIngress",
Match.object_like(
{
"FromPort": 3000,
"ToPort": 3000,
"SourceSecurityGroupId": Match.any_value(),
}
),
)
def test_alb_internet_facing_https_listener():
t = _template()
t.has_resource_properties(
"AWS::ElasticLoadBalancingV2::LoadBalancer", {"Scheme": "internet-facing"}
)
t.has_resource_properties(
"AWS::ElasticLoadBalancingV2::Listener",
Match.object_like(
{"Port": 443, "Protocol": "HTTPS", "Certificates": Match.any_value()}
),
)
def test_no_static_keys_in_stack():
t = _template()
t.resource_count_is("AWS::IAM::User", 0)
t.resource_count_is("AWS::IAM::AccessKey", 0)
def test_instance_role_scoped_and_uses_ssm():
t = _template()
# Session Manager (no SSH) — the SSM managed policy is attached.
t.has_resource_properties(
"AWS::IAM::Role",
Match.object_like(
{
"ManagedPolicyArns": Match.array_with(
[
{
"Fn::Join": [
"",
Match.array_with(
[":iam::aws:policy/AmazonSSMManagedInstanceCore"]
),
]
}
]
)
}
),
)
# The instance role must not be able to write the catalog or run wide Athena.
for policy in t.find_resources("AWS::IAM::Policy").values():
for stmt in policy["Properties"]["PolicyDocument"]["Statement"]:
actions = stmt.get("Action", [])
actions = [actions] if isinstance(actions, str) else actions
for a in actions:
if isinstance(a, str):
assert a not in ("glue:*", "athena:*", "s3:*", "*"), (
f"too broad: {a}"
)
assert not a.startswith("glue:Create"), f"no Glue writes: {a}"
assert not a.startswith("glue:Update"), f"no Glue writes: {a}"
def test_root_volume_gp3_and_retained():
_template().has_resource_properties(
"AWS::EC2::Instance",
Match.object_like(
{
"BlockDeviceMappings": Match.array_with(
[
Match.object_like(
{
"Ebs": Match.object_like(
{
"VolumeType": "gp3",
"DeleteOnTermination": False,
"Encrypted": True,
}
)
}
)
]
)
}
),
)
def test_daily_dlm_backup_enabled():
_template().has_resource_properties(
"AWS::DLM::LifecyclePolicy",
Match.object_like(
{
"State": "ENABLED",
"PolicyDetails": Match.object_like({"ResourceTypes": ["INSTANCE"]}),
}
),
)
def test_route53_alias_for_grafana():
_template().has_resource_properties(
"AWS::Route53::RecordSet",
Match.object_like({"Type": "A", "Name": "grafana.seahaven.com."}),
)

View file

@ -1,16 +1,19 @@
"""Synth-level assertions for the Phase 3 analytics dataset.
"""Synth-level assertions for the analytics dataset (Phase 3) and Slack surfaces
(Phase 4).
Synthesizes ``apm-wo-analysis-pipeline`` and asserts the Glue table carries
partition projection, the column schema matches what the classifier writes, the
Athena workgroup enforces its result location, and the classifier role has zero
Glue access (projection means it never touches the catalog). No AWS, no Docker:
``aws:cdk:bundling-stacks=[]`` skips asset bundling so this is a fast offline gate.
Synthesizes ``apm-wo-analysis-pipeline`` and asserts: the Glue table carries
partition projection with the classifier's column schema; the Athena workgroup
enforces its result location; the classifier role has zero Glue access; and the
Phase 4 Slack post + interactions Lambdas exist behind an HTTP API with only the
scoped S3/Secrets/SSM permissions. No AWS, no Docker: ``aws:cdk:bundling-stacks=[]``
skips asset bundling so this is a fast offline gate.
Run with the repo venv:
python -m pytest tests/test_pipeline_synth.py -q
"""
import json
import sys
from pathlib import Path
@ -28,8 +31,16 @@ def _s3_path_ending(suffix: str):
return {"Fn::Join": ["", Match.array_with([suffix])]}
def _cdk_context() -> dict:
"""The real cdk.json context (cert ARN, hosted zone, Slack domain) — the
hand-built App below doesn't auto-load it the way `cdk synth` does."""
ctx = json.loads((CDK_DIR / "cdk.json").read_text())["context"]
ctx["aws:cdk:bundling-stacks"] = []
return ctx
def _template() -> Template:
app = cdk.App(context={"aws:cdk:bundling-stacks": []})
app = cdk.App(context=_cdk_context())
stack = PipelineStack(
app,
"apm-wo-analysis-pipeline",
@ -117,3 +128,89 @@ def test_classifier_role_has_no_glue_access():
a for a in actions if isinstance(a, str) and a.startswith("glue:")
]
assert not offending, f"unexpected Glue access: {offending}"
# ----- Phase 4 — Slack surfaces -----
def test_slack_lambdas_exist():
t = _template()
for fn_name, handler in (
("apm-wo-analysis-slack-post", "handler.handler"),
("apm-wo-analysis-slack-interactions", "interactions.handler"),
):
t.has_resource_properties(
"AWS::Lambda::Function",
{
"FunctionName": fn_name,
"Handler": handler,
"Runtime": "python3.12",
"Architectures": ["arm64"],
},
)
def test_classifier_has_dlq():
# Failed async invocations must surface, not silently drop a day's data.
t = _template()
t.resource_count_is("AWS::SQS::Queue", 1)
t.has_resource_properties(
"AWS::Lambda::Function",
Match.object_like(
{
"FunctionName": "apm-wo-analysis-classifier",
"DeadLetterConfig": Match.any_value(),
}
),
)
def test_interactions_stage_is_throttled():
# The public Slack interactions endpoint caps rate/burst.
_template().has_resource_properties(
"AWS::ApiGatewayV2::Stage",
Match.object_like(
{
"DefaultRouteSettings": {
"ThrottlingRateLimit": 10,
"ThrottlingBurstLimit": 20,
}
}
),
)
def test_interactions_api_routes_post_to_slack_endpoint():
t = _template()
t.resource_count_is("AWS::ApiGatewayV2::Api", 1)
t.has_resource_properties(
"AWS::ApiGatewayV2::Route", {"RouteKey": "POST /slack/interactions"}
)
# Custom domain on apm-wo.seahaven.com + a Route53 alias for it.
t.has_resource_properties(
"AWS::ApiGatewayV2::DomainName", {"DomainName": "apm-wo.seahaven.com"}
)
t.has_resource_properties("AWS::Route53::RecordSet", {"Type": "A"})
def test_slack_roles_have_no_broad_or_write_access():
# The Slack Lambdas should only read analytics/, the Slack secret, and the
# dashboard SSM param — never write S3, never s3:*/secretsmanager:* wildcards.
# Scope to the Slack policies by logical ID (the classifier/drop-uploader
# legitimately hold s3:PutObject).
t = _template()
slack_policies = {
lid: p
for lid, p in t.find_resources("AWS::IAM::Policy").items()
if lid.startswith(("SlackPost", "SlackInteractions"))
}
assert slack_policies, "expected scoped policies for the Slack roles"
for policy in slack_policies.values():
for stmt in policy["Properties"]["PolicyDocument"]["Statement"]:
actions = stmt.get("Action", [])
actions = [actions] if isinstance(actions, str) else actions
for a in actions:
if not isinstance(a, str):
continue
assert a not in ("s3:*", "secretsmanager:*", "*"), f"too broad: {a}"
assert a != "s3:PutObject", "Slack roles must not write S3"