procurement-ingest/cdk/common.py
Adam Moussa f8eb18f02b
Some checks are pending
Deploy / deploy (push) Waiting to run
feat: collapse duplicated CDK into cdk/common.py plain helpers (refactor phase 4) (#112)
The ~379 lines po_stack.py and wo_stack.py defined identically (DynamoDB
alarms, the sender-auth-rejected metric filter + alarm, the standard
per-Lambda alarm set, the Bedrock InvokeModel grant, the raw-email
bucket, the async DLQ, the template-fallback-rate math alarm) move into
cdk/common.py.

Every helper is a PLAIN function taking (scope, id, ...), called with each
stack's own Stack as scope and the exact literal construct ids used inline
before, so every synthesized logical ID is byte-stable. A Construct
subclass would reparent the tree and make CloudFormation attempt to
replace the RETAIN-protected purchase-orders/WorkOrders tables and named
buckets -- data loss -- so it is forbidden. Per-function alarm variance
(PO p99 vs WO p95 duration, po-web-ui throttles+duration only,
site-extractor no DLQ alarm, workorder-web-ui zero alarms) is preserved
through call-site arguments, not baked into the helpers.

make_bedrock_invoke_statement derives the inference-profile and us-east-1
foundation-model ARNs from Stack.of(scope).account/.region instead of the
hardcoded 328440206208/us-east-1 literals. The environment stays
account-agnostic (region-only), so the account resolves to the
AWS::AccountId pseudo-parameter: the derived ARN resolves at deploy to the
same ARN the literal named in-account (a benign in-place IAM policy
update, never a replacement) and is account-portable rather than pinned to
the frozen management account.

The account= pin evaluated for cdk.Environment was deliberately NOT added:
resolving every account-derived value (bucket names, Lambda::Permission
source account, SNS action ARN) to literals makes CloudFormation flag the
RETAIN email buckets as requiring replacement against the deployed
account-agnostic templates -- a data-loss risk that outranks the pin, which
buys nothing (the resolved values are unchanged).

Also: net-new CfnOutputs for the five Lambda function ARNs and the
owned/consumed table names, exact-pin constructs==10.6.0, and fix the
stale aws-cdk-lib 2.259.0 -> 2.261.0 version comment.

The common.py extraction is zero-cdk-diff on both stacks (byte-stable
logical IDs, no asset/property change); the only deltas versus deployed
are the intended benign Bedrock IAM in-place update and the additive
CfnOutputs. Mandatory GPT-4.1 cross-family review ran on the Bedrock IAM
move; its BLOCK was a verified false positive (it read AWS::AccountId as a
wildcard -- it is a deploy-time-resolved concrete value naming one account
and one inference-profile, region is pinned us-east-1, and the grant is
strictly more least-privilege-correct than the hardcoded literal).
2026-07-20 18:14:57 +00:00

331 lines
12 KiB
Python

"""Shared CDK helpers for the procurement-ingest stacks (Phase 4).
Plain free functions extracted from po_stack.py / wo_stack.py. Each takes the
same Stack `scope` and the same literal construct `id` the stacks passed inline,
so every synthesized logical ID is byte-stable. NOTHING here is a Construct
subclass: a subclass would insert a tree node and reparent/replace the
RETAIN-protected tables and named buckets.
"""
from aws_cdk import (
Duration,
RemovalPolicy,
Stack,
aws_cloudwatch as cloudwatch,
aws_cloudwatch_actions as cw_actions,
aws_dynamodb as dynamodb,
aws_iam as iam,
aws_logs as logs,
aws_s3 as s3,
aws_sqs as sqs,
)
# --- DynamoDB alarm operations (moved verbatim from both stacks) ---
_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,
]
# CloudWatch namespace for the log-derived sender-authentication metrics.
_SENDER_AUTH_METRIC_NAMESPACE = "Seahaven/ProcurementIngest"
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))
def add_sender_auth_rejected_alarm(scope, id_prefix, function_name, alarm_topic):
"""Metric-filter + alarm on ``sender_auth_rejected`` warnings (INFRA-107).
A rejected inbound email is skipped without erroring the invocation, so it
is invisible to the Errors/Throttles/DLQ alarms. This turns the structured
warning log into a CloudWatch metric and pages when rejections spike --
catching a silent false-reject storm (allowlist wrong, signing-domain
drift, SES header-format change) that would otherwise discard legitimate
mail while the pipeline reports healthy.
ALARM-only SnsAction to site-alerts; no OK action. The metric filter reads
the function's own log group (imported by the deterministic
``/aws/lambda/<fn>`` name, created by the function's log_retention). A plain
substring pattern is used because Lambda prefixes each line with its own
level/timestamp/request-id, so the JSON payload is not a standalone JSON
log event a `{$.event=...}` pattern could match.
"""
metric_name = f"{function_name}-sender-auth-rejected"
logs.MetricFilter(
scope,
f"{id_prefix}SenderAuthRejectedFilter",
log_group=logs.LogGroup.from_log_group_name(
scope,
f"{id_prefix}LogGroup",
f"/aws/lambda/{function_name}",
),
filter_pattern=logs.FilterPattern.literal('"sender_auth_rejected"'),
metric_namespace=_SENDER_AUTH_METRIC_NAMESPACE,
metric_name=metric_name,
metric_value="1",
default_value=0,
)
cloudwatch.Metric(
namespace=_SENDER_AUTH_METRIC_NAMESPACE,
metric_name=metric_name,
period=Duration.minutes(5),
statistic="Sum",
).create_alarm(
scope,
f"{id_prefix}SenderAuthRejectedAlarm",
alarm_name=f"{function_name}-sender-auth-rejected",
alarm_description=(
f"{function_name} rejected inbound mail on sender authentication "
"(possible allowlist/DKIM-domain drift silently dropping real mail)"
),
threshold=1,
evaluation_periods=6,
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))
def add_standard_lambda_alarms(
scope,
id_prefix,
fn,
name_prefix,
topic,
*,
duration_statistic,
errors=True,
dlq=None,
descriptions,
):
"""Standard per-Lambda alarm set: errors (optional), throttles, dlq
(optional), duration. Every alarm bespoke-described via `descriptions`
(keys: errors/throttles/dlq/duration passed through VERBATIM). Construct
ids are f"{id_prefix}<Kind>Alarm", alarm names f"{name_prefix}-<kind>",
exactly the inline literals. Order of creation is irrelevant to logical IDs
(ids are explicit) so the fixed errors->throttles->dlq->duration order here
reproduces both PO (dlq before duration) and WO (duration before dlq)
templates identically.
"""
if errors:
fn.metric_errors(period=Duration.minutes(5), statistic="Sum").create_alarm(
scope,
f"{id_prefix}ErrorsAlarm",
alarm_name=f"{name_prefix}-errors",
alarm_description=descriptions["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(topic))
fn.metric_throttles(period=Duration.minutes(5), statistic="Sum").create_alarm(
scope,
f"{id_prefix}ThrottlesAlarm",
alarm_name=f"{name_prefix}-throttles",
alarm_description=descriptions["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(topic))
if dlq is not None:
dlq.metric_approximate_number_of_messages_visible(
period=Duration.minutes(5),
statistic="Maximum",
).create_alarm(
scope,
f"{id_prefix}DlqMessagesAlarm",
alarm_name=f"{name_prefix}-dlq-messages",
alarm_description=descriptions["dlq"],
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(topic))
fn.metric_duration(
period=Duration.minutes(5),
statistic=duration_statistic,
).create_alarm(
scope,
f"{id_prefix}DurationAlarm",
alarm_name=f"{name_prefix}-duration",
alarm_description=descriptions["duration"],
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(topic))
def make_bedrock_invoke_statement(scope):
"""Return the Bedrock InvokeModel PolicyStatement for the email processor.
The us.* inference profile can route cross-region, so the grant covers the
inference-profile ARN plus the per-region foundation-model ARNs
(us-east-1/us-east-2/us-west-2). ARN #1 account and ARN #1/#2 region are
DERIVED from Stack.of(scope).account/.region (not hardcoded 328440206208);
ARN #3/#4 regions (us-east-2/us-west-2) are cross-region reach targets, NOT
the stack's own region, so they stay literal. Caller attaches via
fn.add_to_role_policy(...) on the SAME function -> logical-ID-safe.
"""
stack = Stack.of(scope)
account = stack.account
region = stack.region
return iam.PolicyStatement(
actions=[
"bedrock:InvokeModel",
"bedrock:InvokeModelWithResponseStream",
],
resources=[
f"arn:aws:bedrock:{region}:{account}:inference-profile/us.anthropic.claude-haiku-4-5-20251001-v1:0",
f"arn:aws:bedrock:{region}::foundation-model/anthropic.claude-haiku-4-5-20251001-v1:0",
"arn:aws:bedrock:us-east-2::foundation-model/anthropic.claude-haiku-4-5-20251001-v1:0",
"arn:aws:bedrock:us-west-2::foundation-model/anthropic.claude-haiku-4-5-20251001-v1:0",
],
)
def make_email_bucket(scope, id, name_prefix):
"""Raw-email S3 bucket: BLOCK_ALL public, RETAIN, 90-day expiration.
bucket_name = f"{name_prefix}-{account}" (account derived from the stack),
reproducing f"po-ingest-emails-{self.account}" /
f"workorder-ingest-emails-{self.account}" exactly.
"""
return s3.Bucket(
scope,
id,
bucket_name=f"{name_prefix}-{Stack.of(scope).account}",
block_public_access=s3.BlockPublicAccess.BLOCK_ALL,
removal_policy=RemovalPolicy.RETAIN,
lifecycle_rules=[
s3.LifecycleRule(expiration=Duration.days(90)),
],
)
def make_processor_dlq(scope, id):
"""Async-invoke DLQ: 14-day retention, enforce_ssl. CDK-generated name."""
return sqs.Queue(
scope,
id,
retention_period=Duration.days(14),
enforce_ssl=True,
)
def make_fallback_rate_alarm(
scope,
id,
*,
namespace,
alarm_topic,
alarm_name,
alarm_description,
rejected_included,
period,
threshold,
floor,
evaluation_periods,
datapoints_to_alarm,
):
"""Template-fallback-rate MathExpression alarm.
rejected_included=False (PO): numerator FILL(fb,0), denom fb+tmpl.
rejected_included=True (WO): numerator (fb+rej), denom fb+rej+tmpl.
Expression string, FILL, label reproduced BYTE-FOR-BYTE. GREATER_THAN,
NOT_BREACHING. NO element-wise MAX (post-#102). The two DISTINCT rejected
alarms are NOT built here.
"""
fb_metric = cloudwatch.Metric(
namespace=namespace,
metric_name="ParseOutcome",
dimensions_map={"ParseMethod": "ai_fallback"},
statistic="Sum",
period=period,
)
tmpl_metric = cloudwatch.Metric(
namespace=namespace,
metric_name="ParseOutcome",
dimensions_map={"ParseMethod": "template"},
statistic="Sum",
period=period,
)
if rejected_included:
fb_rej_metric = cloudwatch.Metric(
namespace=namespace,
metric_name="ParseOutcome",
dimensions_map={"ParseMethod": "ai_fallback_rejected"},
statistic="Sum",
period=period,
)
using_metrics = {"fb": fb_metric, "rej": fb_rej_metric, "tmpl": tmpl_metric}
sum_terms = "FILL(fb,0)+FILL(rej,0)+FILL(tmpl,0)"
numerator = "(FILL(fb,0)+FILL(rej,0))"
else:
using_metrics = {"fb": fb_metric, "tmpl": tmpl_metric}
sum_terms = "FILL(fb,0)+FILL(tmpl,0)"
numerator = "FILL(fb,0)"
expression = f"IF(({sum_terms})>={floor}, 100*{numerator}/({sum_terms}), 0)"
fallback_rate = cloudwatch.MathExpression(
expression=expression,
using_metrics=using_metrics,
period=period,
label="TemplateFallbackRatePct",
)
fallback_rate.create_alarm(
scope,
id,
alarm_name=alarm_name,
alarm_description=alarm_description,
threshold=threshold,
evaluation_periods=evaluation_periods,
datapoints_to_alarm=datapoints_to_alarm,
comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD,
treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING,
).add_alarm_action(cw_actions.SnsAction(alarm_topic))