diff --git a/cdk/stacks/grafana_stack.py b/cdk/stacks/grafana_stack.py index ee1a713..17ee09d 100644 --- a/cdk/stacks/grafana_stack.py +++ b/cdk/stacks/grafana_stack.py @@ -188,11 +188,12 @@ class GrafanaStack(Stack): block_devices=[ ec2.BlockDevice( device_name="/dev/xvda", - # RETAIN the gp3 root volume (grafana.db lives here). + # 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, ), ) ], diff --git a/cdk/stacks/pipeline_stack.py b/cdk/stacks/pipeline_stack.py index 88765e1..89bbe53 100644 --- a/cdk/stacks/pipeline_stack.py +++ b/cdk/stacks/pipeline_stack.py @@ -56,6 +56,9 @@ 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, ) @@ -225,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", @@ -235,6 +248,7 @@ 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 @@ -383,6 +397,14 @@ 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( diff --git a/lambdas/classifier/handler.py b/lambdas/classifier/handler.py index f678c33..36cc5de 100644 --- a/lambdas/classifier/handler.py +++ b/lambdas/classifier/handler.py @@ -202,19 +202,22 @@ def _build_details(df: pd.DataFrame) -> list[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"]) +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() @@ -222,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 @@ -254,17 +257,33 @@ def handler(event, context): 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." ) - _invoke_slack_post(dt) - 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: diff --git a/tests/test_grafana_synth.py b/tests/test_grafana_synth.py index 42f29b1..b6784c6 100644 --- a/tests/test_grafana_synth.py +++ b/tests/test_grafana_synth.py @@ -147,7 +147,11 @@ def test_root_volume_gp3_and_retained(): Match.object_like( { "Ebs": Match.object_like( - {"VolumeType": "gp3", "DeleteOnTermination": False} + { + "VolumeType": "gp3", + "DeleteOnTermination": False, + "Encrypted": True, + } ) } ) diff --git a/tests/test_pipeline_synth.py b/tests/test_pipeline_synth.py index eec10ac..641201c 100644 --- a/tests/test_pipeline_synth.py +++ b/tests/test_pipeline_synth.py @@ -150,6 +150,36 @@ 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)