"""Pipeline stack: S3, classifier + slack-post Lambdas, Glue, Athena, IAM. 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. 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). """ import os from aws_cdk import ( BundlingOptions, Duration, 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, ) from aws_cdk import ( aws_iam as iam, ) from aws_cdk import ( aws_lambda as lambda_, ) 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, ) from aws_cdk import ( aws_s3_notifications as s3n, ) from aws_cdk import ( aws_secretsmanager as secretsmanager, ) 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" 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 # (lambdas/classifier/handler.py _build_snapshot). Order/names must match the # Parquet columns; `dt` is the projected partition key, not a stored column. SNAPSHOT_COLUMNS = [ ("wo_number", "string"), ("wo_description", "string"), ("equipment_code", "string"), ("site", "string"), ("due_date", "string"), ("department", "string"), ("wo_status", "string"), ("hold_reason", "string"), ("last_comment", "string"), ("last_comment_by", "string"), ("last_comment_date", "string"), ("contractor", "string"), ("contractor_description", "string"), ("category", "string"), ("is_escalation", "boolean"), ("is_action", "boolean"), ("mismatch", "string"), ] _PARQUET_INPUT = "org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat" _PARQUET_OUTPUT = "org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat" _PARQUET_SERDE = "org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe" class PipelineStack(Stack): def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None: super().__init__(scope, construct_id, **kwargs) # Phase 1 — single exports bucket. # Prefixes: raw/ (incoming), analytics/ (per-WO snapshots), athena-results/. self.exports_bucket = s3.Bucket( self, "Exports", bucket_name=f"apm-wo-analysis-exports-{self.account}", encryption=s3.BucketEncryption.S3_MANAGED, block_public_access=s3.BlockPublicAccess.BLOCK_ALL, enforce_ssl=True, removal_policy=RemovalPolicy.RETAIN, lifecycle_rules=[ s3.LifecycleRule( id="expire-raw-exports", prefix="raw/", expiration=Duration.days(90), ), # Athena query output is disposable; don't let it accumulate in # a RETAIN bucket. analytics/ snapshots are kept indefinitely. s3.LifecycleRule( id="expire-athena-results", prefix="athena-results/", expiration=Duration.days(30), ), ], ) # Phase 1 — least-privilege identity for the local drop-folder uploader. # Scoped to s3:PutObject on raw/* only. The access key is created # out-of-band (aws iam create-access-key) and stored in the local # ~/.aws/credentials profile `apm-wo-drop` — never in CloudFormation. self.drop_uploader = iam.User( self, "DropUploader", user_name="apm-wo-drop-uploader" ) self.drop_uploader.add_to_policy( iam.PolicyStatement( sid="PutRawExportsOnly", actions=["s3:PutObject"], resources=[self.exports_bucket.arn_for_objects("raw/*")], ) ) # Phase 2/3 — Glue database for the analytics dataset. glue.CfnDatabase( self, "AnalyticsDb", catalog_id=self.account, database_input=glue.CfnDatabase.DatabaseInputProperty(name=GLUE_DATABASE), ) # Phase 3 — apm_wo_snapshots table over analytics/, partitioned by dt # with partition projection: Athena derives dt from the path, so there # is no crawler, no MSCK REPAIR, and the classifier needs no Glue access. analytics_location = f"s3://{self.exports_bucket.bucket_name}/analytics/" glue.CfnTable( self, "SnapshotsTable", catalog_id=self.account, database_name=GLUE_DATABASE, table_input=glue.CfnTable.TableInputProperty( name=GLUE_TABLE, table_type="EXTERNAL_TABLE", partition_keys=[ glue.CfnTable.ColumnProperty(name="dt", type="string"), ], parameters={ "classification": "parquet", "EXTERNAL": "TRUE", "projection.enabled": "true", "projection.dt.type": "date", "projection.dt.format": "yyyy-MM-dd", "projection.dt.range": "2026-01-01,NOW", "storage.location.template": f"{analytics_location}dt=${{dt}}/", }, storage_descriptor=glue.CfnTable.StorageDescriptorProperty( location=analytics_location, input_format=_PARQUET_INPUT, output_format=_PARQUET_OUTPUT, serde_info=glue.CfnTable.SerdeInfoProperty( serialization_library=_PARQUET_SERDE ), columns=[ glue.CfnTable.ColumnProperty(name=name, type=type_) for name, type_ in SNAPSHOT_COLUMNS ], ), ), ) # Phase 3 — dedicated Athena workgroup, enforced result location + SSE-S3. athena.CfnWorkGroup( self, "Workgroup", name=ATHENA_WORKGROUP, recursive_delete_option=True, work_group_configuration=athena.CfnWorkGroup.WorkGroupConfigurationProperty( enforce_work_group_configuration=True, publish_cloud_watch_metrics_enabled=True, result_configuration=athena.CfnWorkGroup.ResultConfigurationProperty( output_location=f"s3://{self.exports_bucket.bucket_name}/athena-results/", encryption_configuration=athena.CfnWorkGroup.EncryptionConfigurationProperty( encryption_option="SSE_S3", ), ), ), ) # Phase 2 — classifier Lambda, S3-triggered on the raw/ prefix. # Deps (awswrangler/pandas/pyarrow/openpyxl) are Docker-bundled for ARM64 # from lambdas/classifier/requirements.txt; boto3 ships in the runtime. classifier_logs = logs.LogGroup( self, "ClassifierLogs", log_group_name="/aws/lambda/apm-wo-analysis-classifier", retention=logs.RetentionDays.TWO_MONTHS, removal_policy=RemovalPolicy.DESTROY, ) self.classifier_fn = lambda_.Function( self, "Classifier", function_name="apm-wo-analysis-classifier", runtime=lambda_.Runtime.PYTHON_3_12, architecture=lambda_.Architecture.ARM_64, handler="handler.handler", memory_size=512, timeout=Duration.seconds(120), log_group=classifier_logs, environment={"APM_HAIKU_FALLBACK": "on"}, code=lambda_.Code.from_asset( os.path.join(LAMBDAS_DIR, "classifier"), 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", ], ), ), ) # S3 trigger: any .xlsx/.csv landing under raw/ invokes the classifier. for suffix in (".xlsx", ".csv"): self.exports_bucket.add_event_notification( s3.EventType.OBJECT_CREATED, s3n.LambdaDestination(self.classifier_fn), s3.NotificationKeyFilter(prefix="raw/", suffix=suffix), ) # IAM — least privilege: read raw/, read+write analytics/, and read the # Anthropic key for the Haiku fallback. No Glue access: the table is # CDK-defined with partition projection, so the classifier only writes # 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/*") secretsmanager.Secret.from_secret_name_v2( self, "AnthropicKey", ANTHROPIC_SECRET ).grant_read(self.classifier_fn) # ----- 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", ], ), ), ) self.exports_bucket.grant_read(fn, "analytics/*") 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 ), ) # 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, ) ), )