mirror of
https://github.com/Sea-Haven-Industries/apm-wo-analysis.git
synced 2026-10-04 16:02:04 +00:00
Compare commits
6 commits
c04c239774
...
d884de8fcc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d884de8fcc | ||
|
|
dce53fdc86 | ||
|
|
3c8b7704f6 | ||
|
|
68f3acded8 | ||
|
|
ea54cb1e60 | ||
|
|
f193754c27 |
12 changed files with 81 additions and 31 deletions
|
|
@ -25,7 +25,9 @@ REPO
|
|||
dnf install -y grafana
|
||||
|
||||
# --- Athena datasource plugin (pinned for reproducibility) ---
|
||||
grafana-cli --pluginsDir /var/lib/grafana/plugins plugins install grafana-athena-datasource "${PLUGIN_VERSION}"
|
||||
# --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'
|
||||
|
|
@ -50,8 +52,8 @@ 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
|
||||
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/dashboards/" /var/lib/grafana/dashboards/ --delete
|
||||
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
|
||||
|
|
@ -64,8 +66,8 @@ systemctl enable --now grafana-server
|
|||
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
|
||||
aws s3 sync "s3://${CONFIG_BUCKET}/${CONFIG_PREFIX}/dashboards/" /var/lib/grafana/dashboards/ --delete
|
||||
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
|
||||
|
|
|
|||
|
|
@ -12,6 +12,6 @@
|
|||
"grafanaPrivateSubnetIds": ["subnet-04e38c507e96f1926", "subnet-0a0b4fc6f296dfba5"],
|
||||
"grafanaDomain": "grafana.seahaven.com",
|
||||
"officeCidrs": ["47.21.61.4/32", "96.250.164.146/32"],
|
||||
"athenaPluginVersion": "2.18.2"
|
||||
"athenaPluginVersion": "3.2.0"
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -64,6 +64,11 @@ 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"
|
||||
|
|
@ -231,6 +236,15 @@ class PipelineStack(Stack):
|
|||
timeout=Duration.seconds(120),
|
||||
log_group=classifier_logs,
|
||||
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(
|
||||
|
|
@ -260,6 +274,9 @@ 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)
|
||||
|
|
@ -321,7 +338,8 @@ class PipelineStack(Stack):
|
|||
),
|
||||
),
|
||||
)
|
||||
self.exports_bucket.grant_read(fn, "analytics/*")
|
||||
# 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
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 0,
|
||||
"multi": false,
|
||||
"includeAll": false,
|
||||
|
|
@ -57,7 +57,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 1,
|
||||
"multi": true,
|
||||
"includeAll": true,
|
||||
|
|
@ -84,7 +84,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 1,
|
||||
"multi": true,
|
||||
"includeAll": true,
|
||||
|
|
@ -111,7 +111,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 1,
|
||||
"multi": true,
|
||||
"includeAll": true,
|
||||
|
|
@ -138,7 +138,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 1,
|
||||
"multi": true,
|
||||
"includeAll": true,
|
||||
|
|
@ -165,7 +165,7 @@
|
|||
"region": "us-east-1"
|
||||
}
|
||||
},
|
||||
"refresh": 2,
|
||||
"refresh": 1,
|
||||
"sort": 1,
|
||||
"multi": true,
|
||||
"includeAll": true,
|
||||
|
|
@ -179,7 +179,7 @@
|
|||
"panels": [
|
||||
{
|
||||
"id": 1,
|
||||
"type": "bar-chart",
|
||||
"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 },
|
||||
|
|
@ -366,7 +366,7 @@
|
|||
},
|
||||
{
|
||||
"id": 4,
|
||||
"type": "bar-chart",
|
||||
"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 },
|
||||
|
|
|
|||
|
|
@ -1,5 +1,8 @@
|
|||
# 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
|
||||
|
|
@ -7,7 +10,7 @@ datasources:
|
|||
uid: athena
|
||||
isDefault: true
|
||||
jsonData:
|
||||
authType: ec2_iam_role
|
||||
authType: default
|
||||
defaultRegion: us-east-1
|
||||
catalog: AwsDataCatalog
|
||||
database: apm_wo_analysis
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
@ -237,17 +238,19 @@ 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",
|
||||
)
|
||||
# details.json — per-WO index the slack-post Lambda reads for drill-down modals.
|
||||
_s3.put_object(
|
||||
Bucket=bucket,
|
||||
Key=f"{ANALYTICS_PREFIX}/dt={dt}/details.json",
|
||||
Key=f"{META_PREFIX}/dt={dt}/details.json",
|
||||
Body=json.dumps(_build_details(df)).encode("utf-8"),
|
||||
ContentType="application/json",
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -282,8 +282,11 @@ def build_daily_summary(
|
|||
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})", "drill_category", cat)
|
||||
_button(f"{short_label} ({count})", f"drill_category:{cat}", cat)
|
||||
)
|
||||
|
||||
# Dashboard link button always present (no action_id — url button).
|
||||
|
|
|
|||
|
|
@ -25,11 +25,11 @@ def handler(event, context):
|
|||
"""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_analytics_json(dt, "summary.json")
|
||||
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_analytics_json(_yesterday(dt), "summary.json")
|
||||
yesterday = slackio.read_meta_json(_yesterday(dt), "summary.json")
|
||||
|
||||
client = slackio.web_client()
|
||||
channel = slackio.channel_id()
|
||||
|
|
@ -44,7 +44,7 @@ def handler(event, context):
|
|||
# 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_analytics_json(dt, "details.json") or []
|
||||
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:
|
||||
|
|
|
|||
|
|
@ -49,7 +49,10 @@ def handler(event, context):
|
|||
return {"statusCode": 200, "body": ""}
|
||||
|
||||
action = (payload.get("actions") or [{}])[0]
|
||||
field = _FILTER_FIELD.get(action.get("action_id"))
|
||||
# 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:
|
||||
|
|
@ -57,7 +60,7 @@ def handler(event, context):
|
|||
|
||||
# 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_analytics_json(dt, "details.json") or []
|
||||
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(
|
||||
|
|
|
|||
|
|
@ -60,9 +60,10 @@ def verify_signature(body: str, timestamp: str, signature: str) -> bool:
|
|||
return verifier.is_valid(body=body, timestamp=timestamp, signature=signature)
|
||||
|
||||
|
||||
def read_analytics_json(dt: str, name: str):
|
||||
"""Read ``analytics/dt=<dt>/<name>`` as JSON, or None if absent."""
|
||||
key = f"analytics/dt={dt}/{name}"
|
||||
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:
|
||||
|
|
|
|||
|
|
@ -151,6 +151,20 @@ class TestDailySummaryStructure:
|
|||
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")
|
||||
|
|
@ -191,7 +205,8 @@ class TestDailySummaryStructure:
|
|||
if block.get("type") != "actions":
|
||||
continue
|
||||
for elem in block.get("elements", []):
|
||||
if elem.get("action_id") == "drill_category":
|
||||
# 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"
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue