procurement-ingest/cdk/po_stack.py

545 lines
22 KiB
Python

"""CDK stack for the Coupa PO email ingestion pipeline."""
import aws_cdk as cdk
from aws_cdk import (
Duration,
RemovalPolicy,
Stack,
aws_cloudwatch as cloudwatch,
aws_cloudwatch_actions as cw_actions,
aws_dynamodb as dynamodb,
aws_kms as kms,
aws_lambda as lambda_,
aws_lambda_event_sources as lambda_event_sources,
aws_logs as logs,
aws_s3 as s3,
aws_s3_notifications as s3n,
aws_ses as ses,
aws_ses_actions as ses_actions,
aws_secretsmanager as secretsmanager,
aws_sns as sns,
aws_sqs as sqs,
aws_ssm as ssm,
)
from constructs import Construct
# Operations these tables actually issue (PutItem/UpdateItem/DeleteItem writes,
# GetItem/Query/BatchGetItem reads). DynamoDB emits ThrottledRequests/SystemErrors
# keyed by TableName + Operation only, so the CDK *_for_operations helpers (which
# render a SUM MathExpression across these per-operation metrics) are the correct,
# non-deprecated way to roll a table up to a single alarmable series.
_DDB_ALARM_OPERATIONS = [
dynamodb.Operation.GET_ITEM,
dynamodb.Operation.BATCH_GET_ITEM,
dynamodb.Operation.QUERY,
dynamodb.Operation.SCAN,
dynamodb.Operation.PUT_ITEM,
dynamodb.Operation.UPDATE_ITEM,
dynamodb.Operation.DELETE_ITEM,
dynamodb.Operation.BATCH_WRITE_ITEM,
]
def _add_ddb_alarms(scope, id_prefix, table, alarm_name_prefix, alarm_topic):
"""Add throttle + system-error alarms for a DynamoDB table.
Both fire on any non-zero datapoint in a 5-min window. ALARM-only SnsAction
to site-alerts (no OK action); TreatMissingData NOT_BREACHING.
"""
table.metric_throttled_requests_for_operations(
operations=_DDB_ALARM_OPERATIONS,
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
scope,
f"{id_prefix}ThrottlesAlarm",
alarm_name=f"{alarm_name_prefix}-throttles",
alarm_description=f"{alarm_name_prefix} DynamoDB throttled requests",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
table.metric_system_errors_for_operations(
operations=_DDB_ALARM_OPERATIONS,
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
scope,
f"{id_prefix}SystemErrorsAlarm",
alarm_name=f"{alarm_name_prefix}-system-errors",
alarm_description=f"{alarm_name_prefix} DynamoDB server-side (5xx) errors",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
class PoIngestStack(Stack):
def __init__(self, scope: Construct, construct_id: str, **kwargs):
super().__init__(scope, construct_id, **kwargs)
# --- Shared alarm SNS topic (site-alerts) ---
# Imported once near the top so every alarm in this stack reuses the same
# Topic construct instance (avoids duplicate logical IDs). ALARM-only
# SnsAction; no OK action, per the CloudWatch-alarm preference. The
# topic's CMK (alias/seahaven-alarm-topics) lives on the topic itself.
alarm_topic = sns.Topic.from_topic_arn(
self,
"SiteAlertsTopic",
f"arn:aws:sns:{self.region}:{self.account}:site-alerts",
)
# --- S3 bucket for raw emails ---
email_bucket = s3.Bucket(
self,
"EmailBucket",
bucket_name=f"po-ingest-emails-{self.account}",
block_public_access=s3.BlockPublicAccess.BLOCK_ALL,
removal_policy=RemovalPolicy.RETAIN,
lifecycle_rules=[
s3.LifecycleRule(expiration=Duration.days(90)),
],
)
# --- Shared customer-managed CMK for sensitive DynamoDB tables ---
# Owned by the account-baseline app (alias/seahaven-dynamodb, INFRA-95 /
# M-3); ARN published to SSM. The purchase-orders table was migrated to
# SSE-KMS out-of-band, so declaring encryption_key here reconciles the
# drift and — via grant_read_write_data below — propagates the required
# kms:Decrypt/GenerateDataKey/DescribeKey to the consumer roles.
dynamodb_cmk = kms.Key.from_key_arn(
self,
"DynamoDbCmk",
ssm.StringParameter.value_for_string_parameter(
self, "/seahaven/dynamodb/cmk-arn"
),
)
# --- Shared customer-managed CMK for CloudWatch log encryption ---
# alias/seahaven-logs. The po-email-processor log group was already
# associated with this CMK out-of-band; declaring encryption_key on the
# migrated LogGroups below reconciles that drift and hardens the remaining
# groups (INFRA-114, per cross-review). The key policy already permits
# logs.us-east-1.amazonaws.com (verified in use on po-email-processor).
logs_cmk = kms.Key.from_key_arn(
self,
"LogsKey",
"arn:aws:kms:us-east-1:328440206208:key/b748750c-3b26-478d-acfb-d0126cc97f56",
)
# --- Purchase-orders DynamoDB table ---
# Owned by this stack. Streams enabled for the site-extractor pipeline.
# Other stacks (seahaven-slack-bot) reference this table via fromTableName().
po_table = dynamodb.Table(
self,
"PurchaseOrdersTable",
table_name="purchase-orders",
partition_key=dynamodb.Attribute(
name="po_number",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
stream=dynamodb.StreamViewType.NEW_IMAGE,
encryption=dynamodb.TableEncryption.CUSTOMER_MANAGED,
encryption_key=dynamodb_cmk,
)
# --- Secrets Manager for Anthropic API key ---
anthropic_secret = secretsmanager.Secret(
self,
"AnthropicApiKey",
secret_name="po-ingest/anthropic-api-key",
description="Anthropic API key for Coupa PO email parsing",
removal_policy=RemovalPolicy.RETAIN,
)
# --- DLQ for failed async invocations (INFRA-41 / audit H-8) ---
# SES → S3 → Lambda is async; without an OnFailure destination a failed
# parse (bad email, transient error) is silently dropped after Lambda's
# retries. CDK generates the queue name to avoid colliding with the
# interim CLI-created po-email-processor-dlq (removed post-deploy).
email_processor_dlq = sqs.Queue(
self,
"EmailProcessorDlq",
retention_period=Duration.days(14),
enforce_ssl=True,
)
# --- Explicit log group (INFRA-114) ---
# Replaces the deprecated log_retention prop, which provisioned a
# LogRetention custom resource whose role held logs:PutRetentionPolicy/
# DeleteRetentionPolicy on Resource "*" (checkov CKV_AWS_111). The name
# matches Lambda's default (/aws/lambda/<function-name>) so the function
# keeps writing to the same group; RETAIN matches the repo convention for
# stateful resources.
email_processor_log_group = logs.LogGroup(
self,
"EmailProcessorLogGroup",
log_group_name="/aws/lambda/po-email-processor",
retention=logs.RetentionDays.TWO_MONTHS,
encryption_key=logs_cmk,
removal_policy=RemovalPolicy.RETAIN,
)
# --- Lambda function ---
email_processor = lambda_.Function(
self,
"EmailProcessor",
function_name="po-email-processor",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset(
"../lambdas/po/email_processor",
bundling=cdk.BundlingOptions(
image=lambda_.Runtime.PYTHON_3_12.bundling_image,
command=[
"bash",
"-c",
"pip install --platform manylinux2014_aarch64 --only-binary=:all: "
"-r requirements.txt -t /asset-output && "
"cp handler.py /asset-output/",
],
),
),
timeout=Duration.seconds(60),
memory_size=256,
log_group=email_processor_log_group,
dead_letter_queue=email_processor_dlq,
environment={
"PO_TABLE": "purchase-orders",
"ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn,
},
)
# Grant permissions
email_bucket.grant_read(email_processor)
po_table.grant_read_write_data(email_processor)
anthropic_secret.grant_read(email_processor)
# --- Errors alarm (INFRA-41 / audit H-8) ---
# ALARM-only (no OK action, per the CloudWatch-alarm preference) to the
# shared site-alerts topic. Any errored invocation in a 5-min window pages.
email_processor.metric_errors(
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
self,
"EmailProcessorErrorsAlarm",
alarm_name="po-email-processor-errors",
alarm_description="po-email-processor async invocation errors",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- Throttles alarm: po-email-processor ---
# Any throttled invocation (concurrency cap hit) in a 5-min window pages.
# ALARM-only to site-alerts; no OK action; NOT_BREACHING when no data.
email_processor.metric_throttles(
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
self,
"EmailProcessorThrottlesAlarm",
alarm_name="po-email-processor-throttles",
alarm_description="po-email-processor invocation throttles",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- DLQ messages-present alarm ---
# Pages when any message lands in the EmailProcessorDlq: a message here
# means a PO email was permanently dropped after Lambda exhausted its
# async retries. Maximum over a single 5-min window > 0 fires; missing
# data (no messages metric emitted) is not breaching. Reuses the shared
# site-alerts topic, ALARM-only, like the errors alarm above. The metric
# helper derives the QueueName dimension from the queue construct, so the
# alarm tracks the CDK-generated queue name without hardcoding it.
email_processor_dlq.metric_approximate_number_of_messages_visible(
period=Duration.minutes(5),
statistic="Maximum",
).create_alarm(
self,
"EmailProcessorDlqMessagesAlarm",
alarm_name="po-email-processor-dlq-messages",
alarm_description="po-email-processor DLQ has messages (dropped PO emails)",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- Duration alarm: po-email-processor (orphan adoption) ---
# Adopts the orphaned CLI alarm Lambda-Duration-po-email-processor under
# the repo's <fn>-duration naming (NEW logical name → no deploy collision;
# delete the orphan post-deploy). p99 / 45000 ms
# (75% of the 60s timeout) / eval 3 of 3 — tighter than the orphan's
# Maximum>=48000 / 1-of-1.
email_processor.metric_duration(
period=Duration.minutes(5),
statistic="p99",
).create_alarm(
self,
"EmailProcessorDurationAlarm",
alarm_name="po-email-processor-duration",
alarm_description="po-email-processor p99 duration approaching the 60s timeout",
threshold=45000,
evaluation_periods=3,
datapoints_to_alarm=2,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# S3 event notification → Lambda
email_bucket.add_event_notification(
s3.EventType.OBJECT_CREATED,
s3n.LambdaDestination(email_processor),
s3.NotificationKeyFilter(prefix="inbound/"),
)
# --- SES Receipt Rule ---
# Reuse the existing INBOUND_MAIL rule set (shared with workorder-ingest)
rule_set = ses.ReceiptRuleSet.from_receipt_rule_set_name(
self,
"ExistingRuleSet",
"INBOUND_MAIL",
)
rule_set.add_rule(
"PoEmailRule",
recipients=["amazon_po@int.seahaven.com"],
actions=[
ses_actions.S3(
bucket=email_bucket,
object_key_prefix="inbound/",
),
],
)
# --- Web UI Lambda ---
web_ui_log_group = logs.LogGroup(
self,
"WebUILogGroup",
log_group_name="/aws/lambda/po-web-ui",
retention=logs.RetentionDays.TWO_MONTHS,
encryption_key=logs_cmk,
removal_policy=RemovalPolicy.RETAIN,
)
web_ui = lambda_.Function(
self,
"WebUI",
function_name="po-web-ui",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/po/web_ui"),
timeout=Duration.seconds(60),
memory_size=256,
log_group=web_ui_log_group,
environment={
"PO_TABLE": "purchase-orders",
},
)
po_table.grant_read_data(web_ui)
# --- Throttles alarm: po-web-ui ---
web_ui.metric_throttles(
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
self,
"WebUiThrottlesAlarm",
alarm_name="po-web-ui-throttles",
alarm_description="po-web-ui invocation throttles",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- Duration alarm: po-web-ui ---
# Net-new (no orphan exists for this function).
# p99 / 45000 ms (75% of the 60s timeout) / eval 3, datapoints 2.
web_ui.metric_duration(
period=Duration.minutes(5),
statistic="p99",
).create_alarm(
self,
"WebUiDurationAlarm",
alarm_name="po-web-ui-duration",
alarm_description="po-web-ui p99 duration approaching the 60s timeout",
threshold=45000,
evaluation_periods=3,
datapoints_to_alarm=2,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# Public Function URL removed 2026-06-08 (INFRA-74 / audit C-5): the
# unauthenticated FunctionUrlAuthType.NONE URL was deleted out-of-band
# via CLI. Removing the construct (and its auto-generated Principal:*
# invoke permission) reconciles IaC with the live state.
# --- Verified sites table (extracted from PO ship-to addresses) ---
verified_sites_table = dynamodb.Table(
self,
"VerifiedSitesTable",
table_name="verified-sites",
partition_key=dynamodb.Attribute(
name="siteCode",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
# by-state GSI removed 2026-06-03 (audit M-20): 0 reads in 30d against
# 518 WCU of write amplification. Re-add if a state-level query path ships.
# --- Site extractor Lambda (DynamoDB Streams → verified-sites) ---
site_extractor_log_group = logs.LogGroup(
self,
"SiteExtractorLogGroup",
log_group_name="/aws/lambda/po-ingest-site-extractor",
retention=logs.RetentionDays.TWO_MONTHS,
encryption_key=logs_cmk,
removal_policy=RemovalPolicy.RETAIN,
)
site_extractor = lambda_.Function(
self,
"SiteExtractor",
function_name="po-ingest-site-extractor",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
code=lambda_.Code.from_asset("../lambdas/po/site_extractor"),
timeout=Duration.seconds(60),
memory_size=256,
log_group=site_extractor_log_group,
environment={
"VERIFIED_SITES_TABLE": verified_sites_table.table_name,
"PENDING_REVIEW_TABLE": "pending-site-review",
},
)
verified_sites_table.grant_read_write_data(site_extractor)
site_extractor.add_event_source(
lambda_event_sources.DynamoEventSource(
po_table,
starting_position=lambda_.StartingPosition.TRIM_HORIZON,
batch_size=10,
max_batching_window=Duration.seconds(30),
bisect_batch_on_error=True,
retry_attempts=3,
)
)
# --- Errors alarm: po-ingest-site-extractor ---
# Stream-consumer errors retry per the event-source config, but a
# persistent failure stalls the verified-sites pipeline. ALARM-only to
# site-alerts; no OK action; NOT_BREACHING when no data.
site_extractor.metric_errors(
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
self,
"SiteExtractorErrorsAlarm",
alarm_name="po-ingest-site-extractor-errors",
alarm_description="po-ingest-site-extractor invocation errors",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- Throttles alarm: po-ingest-site-extractor ---
site_extractor.metric_throttles(
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
self,
"SiteExtractorThrottlesAlarm",
alarm_name="po-ingest-site-extractor-throttles",
alarm_description="po-ingest-site-extractor invocation throttles",
threshold=0,
evaluation_periods=1,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
# --- Duration alarm: po-ingest-site-extractor ---
# Net-new (no orphan exists for this function).
# p99 / 45000 ms (75% of the 60s timeout) / eval 3, datapoints 2.
site_extractor.metric_duration(
period=Duration.minutes(5),
statistic="p99",
).create_alarm(
self,
"SiteExtractorDurationAlarm",
alarm_name="po-ingest-site-extractor-duration",
alarm_description="po-ingest-site-extractor p99 duration approaching the 60s timeout",
threshold=45000,
evaluation_periods=3,
datapoints_to_alarm=2,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))
cdk.CfnOutput(
self,
"VerifiedSitesTableName",
value=verified_sites_table.table_name,
description="Verified site addresses extracted from POs",
)
# --- Pending site review table (POs with no extractable site code) ---
pending_review_table = dynamodb.Table(
self,
"PendingSiteReviewTable",
table_name="pending-site-review",
partition_key=dynamodb.Attribute(
name="po_number",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
removal_policy=RemovalPolicy.RETAIN,
)
pending_review_table.grant_read_write_data(site_extractor)
verified_sites_table.grant_read_data(site_extractor)
# --- DynamoDB throttle + system-error alarms ---
# ThrottledRequests / SystemErrors emit at TableName + Operation only
# (verified against live CloudWatch: no TableName-only rollup exists, and
# metric_throttled_requests is deprecated/invalid in aws-cdk-lib 2.259.0).
# Each table currently has zero throttle/error datapoints, so the series
# only materialise on first occurrence — NOT_BREACHING keeps them OK until
# then.
_add_ddb_alarms(
self, "PurchaseOrdersTable", po_table, "purchase-orders", alarm_topic
)
_add_ddb_alarms(
self,
"VerifiedSitesTable",
verified_sites_table,
"verified-sites",
alarm_topic,
)
_add_ddb_alarms(
self,
"PendingSiteReviewTable",
pending_review_table,
"pending-site-review",
alarm_topic,
)