mirror of
https://github.com/Sea-Haven-Industries/apm-wo-analysis.git
synced 2026-09-30 11:13:14 +00:00
Implement the core classification engine and wire it into the pipeline stack. classify.py: two-axis classifier — HTML-strip, comment-intent regex buckets (escalations → status inquiry, most-specific first), Hold Reason / WO Status structured state, comment-vs-state mismatch detector, and a Claude Haiku fallback (Secrets Manager key) reserved for ambiguous free-text. Exports ESCALATION_CATEGORIES / ACTION_NEEDED_CATEGORIES. handler.py: S3-triggered handler — parse xlsx/csv, classify each non-blank row, write a per-WO Parquet snapshot to analytics/dt=YYYY-MM-DD/ (registers the Glue partition via awswrangler) and a summary.json for slack-post (Phase 4). pipeline_stack.py: Glue database, ARM64 Python 3.12 classifier Lambda (Docker-bundled deps), S3 raw/ notification (.xlsx/.csv), and least-privilege IAM (read raw/, read-write analytics/, scoped Glue catalog, read Anthropic key). Smoke-tested against the real export: 347 rows, "Other" at 5.2% (target ~9%), 18 mismatches flagged. 7/7 unit + smoke tests pass; cdk synth green.
170 lines
6.3 KiB
Python
170 lines
6.3 KiB
Python
"""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 4 — slack-post Lambda + scoped IAM (TODO).
|
|
|
|
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_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_s3 as s3,
|
|
)
|
|
from aws_cdk import (
|
|
aws_s3_notifications as s3n,
|
|
)
|
|
from aws_cdk import (
|
|
aws_secretsmanager as secretsmanager,
|
|
)
|
|
from constructs import Construct
|
|
|
|
GLUE_DATABASE = "apm_wo_analysis"
|
|
GLUE_TABLE = "apm_wo_snapshots"
|
|
ANTHROPIC_SECRET = "apm-wo-analysis/anthropic-api-key"
|
|
LAMBDAS_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "lambdas")
|
|
|
|
|
|
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),
|
|
)
|
|
],
|
|
)
|
|
|
|
# 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. The classifier
|
|
# registers the apm_wo_snapshots table/partitions into it via awswrangler;
|
|
# Phase 3 swaps in the partition-projection table definition + Athena.
|
|
glue.CfnDatabase(
|
|
self,
|
|
"AnalyticsDb",
|
|
catalog_id=self.account,
|
|
database_input=glue.CfnDatabase.DatabaseInputProperty(name=GLUE_DATABASE),
|
|
)
|
|
|
|
# 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/, catalog the
|
|
# snapshots table, and read the Anthropic key for the Haiku fallback.
|
|
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)
|
|
|
|
# Phase 4 — slack-post Lambda + scoped IAM; classifier async-invokes it. TODO
|