diff --git a/README.md b/README.md index 7f094ab..392df6e 100644 --- a/README.md +++ b/README.md @@ -142,7 +142,7 @@ cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role ## Status -**Phase 2 — classifier (in review).** Build-out proceeds per +**Phase 3 — analytics dataset (in review).** Build-out proceeds per [`docs/BUILD.md`](./docs/BUILD.md): ingestion → classifier → Glue/Athena → Slack → Grafana → docs. @@ -151,5 +151,8 @@ cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role deployed; PR #7 open. - **Phase 2** classifier Lambda + Glue database + S3 `raw/` trigger — implemented and `cdk synth`-green; PR #8 open (stacked on Phase 1, not yet deployed). -- **Phases 3–6** (Glue/Athena query layer, Slack, Grafana, final docs) — not - started. +- **Phase 3** `apm_wo_snapshots` projection table + Athena workgroup — + implemented and `cdk synth`-green; PR #9 open (stacked on Phase 2). The + classifier writes pure Parquet and holds no Glue access (projection handles + partitions). +- **Phases 4–6** (Slack, Grafana, final docs) — not started. diff --git a/cdk/stacks/pipeline_stack.py b/cdk/stacks/pipeline_stack.py index cba4334..a258317 100644 --- a/cdk/stacks/pipeline_stack.py +++ b/cdk/stacks/pipeline_stack.py @@ -1,8 +1,8 @@ """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 (this file). -Phase 3 — Athena workgroup + partition-projection table (TODO). +Phase 2 — Glue database + S3-triggered classifier Lambda. +Phase 3 — partition-projection Glue table + Athena workgroup (this file). Phase 4 — slack-post Lambda + scoped IAM (TODO). Lambdas: Python 3.12, ARM64, explicit LogGroup with 60-day retention. @@ -17,6 +17,9 @@ from aws_cdk import ( RemovalPolicy, Stack, ) +from aws_cdk import ( + aws_athena as athena, +) from aws_cdk import ( aws_glue as glue, ) @@ -42,9 +45,37 @@ 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" 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: @@ -65,7 +96,14 @@ class PipelineStack(Stack): 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), + ), ], ) @@ -84,9 +122,7 @@ class PipelineStack(Stack): ) ) - # Phase 2/3 — Glue database for the analytics dataset. The classifier - # registers the apm_wo_snapshots table/partitions into it via awswrangler; - # Phase 3 swaps in the partition-projection table definition + Athena. + # Phase 2/3 — Glue database for the analytics dataset. glue.CfnDatabase( self, "AnalyticsDb", @@ -94,6 +130,63 @@ class PipelineStack(Stack): 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. @@ -138,31 +231,12 @@ class PipelineStack(Stack): s3.NotificationKeyFilter(prefix="raw/", suffix=suffix), ) - # IAM — least privilege: read raw/, read+write analytics/, catalog the - # snapshots table, and read the Anthropic key for the Haiku fallback. + # 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/*") - self.classifier_fn.add_to_role_policy( - iam.PolicyStatement( - sid="GlueCatalogSnapshots", - actions=[ - "glue:GetDatabase", - "glue:GetTable", - "glue:CreateTable", - "glue:UpdateTable", - "glue:GetPartition", - "glue:GetPartitions", - "glue:CreatePartition", - "glue:BatchCreatePartition", - "glue:UpdatePartition", - ], - resources=[ - f"arn:aws:glue:{self.region}:{self.account}:catalog", - f"arn:aws:glue:{self.region}:{self.account}:database/{GLUE_DATABASE}", - f"arn:aws:glue:{self.region}:{self.account}:table/{GLUE_DATABASE}/*", - ], - ) - ) secretsmanager.Secret.from_secret_name_v2( self, "AnthropicKey", ANTHROPIC_SECRET ).grant_read(self.classifier_fn) diff --git a/lambdas/classifier/handler.py b/lambdas/classifier/handler.py index 1c5999a..ee8f32c 100644 --- a/lambdas/classifier/handler.py +++ b/lambdas/classifier/handler.py @@ -4,8 +4,9 @@ Triggered on ``s3:ObjectCreated`` under the ``raw/`` prefix regardless of how th file arrives (direct upload or local drop-folder). For each non-blank-comment row it strips HTML, runs the two-axis classifier, derives ``is_escalation`` / ``is_action`` / ``mismatch``, writes a per-WO snapshot to ``analytics/dt=YYYY-MM-DD/`` -as Parquet (registering the Glue partition), and emits a small ``summary.json`` -for the slack-post Lambda to read cheaply. +as Parquet, and emits a small ``summary.json`` for the slack-post Lambda to read +cheaply. The ``apm_wo_snapshots`` Glue table is CDK-defined with partition +projection, so the write needs no Glue catalog access. The slack-post invocation is wired in Phase 4 (the function does not exist yet). Runtime: Python 3.12, ARM64, 512 MB, 120 s. pandas/pyarrow/awswrangler are @@ -29,8 +30,6 @@ import pandas as pd import classify as clf -GLUE_DATABASE = "apm_wo_analysis" -GLUE_TABLE = "apm_wo_snapshots" ANALYTICS_PREFIX = "analytics" # Map snapshot field -> substring matched (case-insensitively) against the export @@ -202,14 +201,16 @@ def handler(event, context): return {"classified": 0, "blank": blank} df["dt"] = dt + # Pure Parquet write — no Glue registration. The apm_wo_snapshots table is + # CDK-defined with partition projection (Phase 3), so Athena derives the dt + # partition from the path and the classifier needs no Glue catalog access. + # overwrite_partitions keeps a same-day re-upload idempotent (replaces dt=). wr.s3.to_parquet( df=df, path=f"s3://{bucket}/{ANALYTICS_PREFIX}/", dataset=True, partition_cols=["dt"], mode="overwrite_partitions", - database=GLUE_DATABASE, - table=GLUE_TABLE, ) summary = _build_summary(df, dt, key, blank) diff --git a/tests/test_pipeline_synth.py b/tests/test_pipeline_synth.py new file mode 100644 index 0000000..19540c9 --- /dev/null +++ b/tests/test_pipeline_synth.py @@ -0,0 +1,119 @@ +"""Synth-level assertions for the Phase 3 analytics dataset. + +Synthesizes ``apm-wo-analysis-pipeline`` and asserts the Glue table carries +partition projection, the column schema matches what the classifier writes, the +Athena workgroup enforces its result location, and the classifier role has zero +Glue access (projection means it never touches the catalog). No AWS, no Docker: +``aws:cdk:bundling-stacks=[]`` skips asset bundling so this is a fast offline gate. + +Run with the repo venv: + + python -m pytest tests/test_pipeline_synth.py -q +""" + +import sys +from pathlib import Path + +import aws_cdk as cdk +from aws_cdk.assertions import Match, Template + +CDK_DIR = Path(__file__).resolve().parents[1] / "cdk" +sys.path.insert(0, str(CDK_DIR)) + +from stacks.pipeline_stack import PipelineStack # noqa: E402 + + +def _s3_path_ending(suffix: str): + """Match an ``Fn::Join`` S3 path (bucket name is a Ref) ending in ``suffix``.""" + return {"Fn::Join": ["", Match.array_with([suffix])]} + + +def _template() -> Template: + app = cdk.App(context={"aws:cdk:bundling-stacks": []}) + stack = PipelineStack( + app, + "apm-wo-analysis-pipeline", + env=cdk.Environment(account="328440206208", region="us-east-1"), + ) + return Template.from_stack(stack) + + +def test_snapshots_table_uses_partition_projection(): + _template().has_resource_properties( + "AWS::Glue::Table", + { + "TableInput": { + "Name": "apm_wo_snapshots", + "PartitionKeys": [{"Name": "dt", "Type": "string"}], + "Parameters": Match.object_like( + { + "projection.enabled": "true", + "projection.dt.type": "date", + "projection.dt.format": "yyyy-MM-dd", + "projection.dt.range": "2026-01-01,NOW", + "storage.location.template": _s3_path_ending( + "/analytics/dt=${dt}/" + ), + } + ), + } + }, + ) + + +def test_snapshots_table_column_schema_matches_classifier(): + # Booleans typed correctly and the columns the BUILD.md prose omitted + # (contractor_description) are present — schema mirrors the writer. + _template().has_resource_properties( + "AWS::Glue::Table", + { + "TableInput": { + "StorageDescriptor": Match.object_like( + { + # array_with is order-sensitive: list patterns in the + # same order the classifier writes them. + "Columns": Match.array_with( + [ + {"Name": "contractor_description", "Type": "string"}, + {"Name": "is_escalation", "Type": "boolean"}, + {"Name": "is_action", "Type": "boolean"}, + {"Name": "mismatch", "Type": "string"}, + ] + ), + } + ), + } + }, + ) + + +def test_athena_workgroup_enforces_results_location(): + _template().has_resource_properties( + "AWS::Athena::WorkGroup", + { + "Name": "apm-wo-analysis", + "WorkGroupConfiguration": Match.object_like( + { + "EnforceWorkGroupConfiguration": True, + "ResultConfiguration": { + "OutputLocation": _s3_path_ending("/athena-results/"), + "EncryptionConfiguration": {"EncryptionOption": "SSE_S3"}, + }, + } + ), + }, + ) + + +def test_classifier_role_has_no_glue_access(): + # Partition projection => the classifier only writes Parquet to S3. No IAM + # policy in the stack should grant any glue:* action. + policies = _template().find_resources("AWS::IAM::Policy") + for policy in policies.values(): + for stmt in policy["Properties"]["PolicyDocument"]["Statement"]: + actions = stmt.get("Action", []) + actions = [actions] if isinstance(actions, str) else actions + offending = [ + a for a in actions if isinstance(a, str) and a.startswith("glue:") + ] + assert not offending, f"unexpected Glue access: {offending}"