Compare commits

..

No commits in common. "eae4d67e0037c2159c21b984506152abfa048c96" and "d884de8fcc3323d696700b8a676b8b962bd1cf48" have entirely different histories.

6 changed files with 37 additions and 108 deletions

View file

@ -188,12 +188,11 @@ class GrafanaStack(Stack):
block_devices=[
ec2.BlockDevice(
device_name="/dev/xvda",
# RETAIN the gp3 root volume (grafana.db lives here), encrypted.
# RETAIN the gp3 root volume (grafana.db lives here).
volume=ec2.BlockDeviceVolume.ebs(
20,
volume_type=ec2.EbsDeviceVolumeType.GP3,
delete_on_termination=False,
encrypted=True,
),
)
],

View file

@ -56,9 +56,6 @@ 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,
)
@ -228,16 +225,6 @@ 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",
@ -248,7 +235,6 @@ 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
@ -397,14 +383,6 @@ class PipelineStack(Stack):
"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(

View file

@ -23,7 +23,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT dt FROM apm_wo_analysis.apm_wo_snapshots ORDER BY dt DESC",
"rawSql": "SELECT DISTINCT dt FROM apm_wo_analysis.apm_wo_snapshots ORDER BY dt DESC",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
@ -49,7 +49,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT site FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND site<>'' ORDER BY 1",
"rawSql": "SELECT DISTINCT site FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND site<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
@ -61,6 +61,7 @@
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": "All",
"current": {},
"options": [],
"hide": 0
@ -75,7 +76,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT department FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND department<>'' ORDER BY 1",
"rawSql": "SELECT DISTINCT department FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND department<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
@ -87,6 +88,7 @@
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": "All",
"current": {},
"options": [],
"hide": 0
@ -101,7 +103,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT category FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND category<>'' ORDER BY 1",
"rawSql": "SELECT DISTINCT category FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND category<>'' ORDER BY 1",
"format": "table",
"connectionArgs": {
"catalog": "AwsDataCatalog",
@ -113,6 +115,7 @@
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": "All",
"current": {},
"options": [],
"hide": 0
@ -127,7 +130,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT wo_status FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND wo_status<>'' ORDER BY 1",
"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",
@ -139,6 +142,7 @@
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": "All",
"current": {},
"options": [],
"hide": 0
@ -153,7 +157,7 @@
"uid": "athena"
},
"query": {
"rawSQL": "SELECT DISTINCT hold_reason FROM apm_wo_analysis.apm_wo_snapshots WHERE dt='$dt' AND hold_reason<>'' ORDER BY 1",
"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",
@ -165,6 +169,7 @@
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": "All",
"current": {},
"options": [],
"hide": 0
@ -189,7 +194,7 @@
"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",
"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",
@ -252,7 +257,7 @@
"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",
"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",
@ -320,7 +325,7 @@
"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",
"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",
@ -376,7 +381,7 @@
"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",
"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",
@ -440,7 +445,7 @@
"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",
"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",
@ -521,7 +526,7 @@
"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",
"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",
@ -676,7 +681,7 @@
"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",
"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",

View file

@ -202,22 +202,19 @@ def _build_details(df: pd.DataFrame) -> list[dict]:
]
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 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"])
def _process_object(bucket: str, key: str, dt: str) -> dict | None:
"""Classify one export into the dt partition + write its summary/details.
Returns the summary dict, or None if the object isn't a usable export."""
if not key.startswith("raw/") or not key.lower().endswith((".xlsx", ".csv")):
print(f"Skipping non-export object: s3://{bucket}/{key}")
return None
return {"skipped": key}
dt = datetime.now(timezone.utc).strftime("%Y-%m-%d")
print(f"Classifying s3://{bucket}/{key} into dt={dt}")
with tempfile.NamedTemporaryFile(suffix=os.path.splitext(key)[1]) as tmp:
_s3.download_fileobj(bucket, key, tmp)
tmp.flush()
@ -225,8 +222,8 @@ def _process_object(bucket: str, key: str, dt: str) -> dict | None:
df, blank = _build_snapshot(header, data)
if df.empty:
print(f"No classifiable rows in {key} (all comments blank); nothing written.")
return None
print("No classifiable rows (all comments blank); nothing written.")
return {"classified": 0, "blank": blank}
df["dt"] = dt
# Pure Parquet write — no Glue registration. The apm_wo_snapshots table is
@ -257,33 +254,17 @@ def _process_object(bucket: str, key: str, dt: str) -> dict | None:
Body=json.dumps(_build_details(df)).encode("utf-8"),
ContentType="application/json",
)
print(
f"Wrote {len(df)} rows, {summary['escalation_total']} escalations "
f"({summary['third_escalation_count']} 3rd), {len(summary['mismatches'])} mismatches."
)
return summary
def handler(event, context):
"""Classify EVERY export in the S3 event (S3 can batch multiple records),
then trigger the Slack post once per affected day. Raises on any failure so
the event is retried / lands in the DLQ rather than being silently dropped."""
processed: list[dict] = []
dts: set[str] = set()
for record in event.get("Records", []):
bucket = record["s3"]["bucket"]["name"]
key = unquote_plus(record["s3"]["object"]["key"])
dt = _event_dt(record)
summary = _process_object(bucket, key, dt)
if summary is not None:
processed.append(
{"key": key, "dt": dt, "classified": summary["classified_total"]}
)
dts.add(dt)
for dt in sorted(dts):
_invoke_slack_post(dt)
return {"processed": processed}
_invoke_slack_post(dt)
return {
"classified": int(len(df)),
"dt": dt,
"summary_key": f"dt={dt}/summary.json",
}
def _invoke_slack_post(dt: str) -> None:

View file

@ -147,11 +147,7 @@ def test_root_volume_gp3_and_retained():
Match.object_like(
{
"Ebs": Match.object_like(
{
"VolumeType": "gp3",
"DeleteOnTermination": False,
"Encrypted": True,
}
{"VolumeType": "gp3", "DeleteOnTermination": False}
)
}
)

View file

@ -150,36 +150,6 @@ def test_slack_lambdas_exist():
)
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)