From eae4d67e0037c2159c21b984506152abfa048c96 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri, 29 May 2026 11:29:14 -0400 Subject: [PATCH] Apply cross-review findings (Phase 2/4/5 hardening) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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). --- cdk/stacks/grafana_stack.py | 3 +- cdk/stacks/pipeline_stack.py | 22 +++++++++++++++ lambdas/classifier/handler.py | 53 ++++++++++++++++++++++++----------- tests/test_grafana_synth.py | 6 +++- tests/test_pipeline_synth.py | 30 ++++++++++++++++++++ 5 files changed, 95 insertions(+), 19 deletions(-) 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)