mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 10:43:14 +00:00
357 lines
14 KiB
Python
357 lines
14 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 make_function_log_group(scope, id_prefix, function_name):
|
|
"""Explicit log group for a Lambda, replacing the deprecated
|
|
``log_retention`` prop (INFRA-114).
|
|
|
|
Making the group a real stack resource (a) drops the LogRetention custom
|
|
resource whose role carried wildcard ``logs:PutRetentionPolicy`` (checkov
|
|
CKV_AWS_111), and (b) makes the group an orderable CFN dependency:
|
|
consumers like the sender-auth metric filter now deploy AFTER the group
|
|
exists. The previous by-name import raced group creation in a fresh
|
|
account and failed the first seahaven-prod deploy (the mgmt account
|
|
masked this because its groups predated the filter).
|
|
|
|
RETAIN matches the repo convention for stateful resources and mirrors the
|
|
old behavior (LogRetention never deleted groups on stack delete). NOTE:
|
|
never ``cdk deploy`` this app to mgmt (328440206208) — PLAT-67 deleted the
|
|
mgmt stacks and left RETAIN cold-archive log groups; a redeploy would
|
|
collide with those leftovers.
|
|
"""
|
|
return logs.LogGroup(
|
|
scope,
|
|
f"{id_prefix}LogGroup",
|
|
log_group_name=f"/aws/lambda/{function_name}",
|
|
retention=logs.RetentionDays.TWO_MONTHS,
|
|
removal_policy=RemovalPolicy.RETAIN,
|
|
)
|
|
|
|
|
|
def add_sender_auth_rejected_alarm(
|
|
scope, id_prefix, function_name, alarm_topic, log_group
|
|
):
|
|
"""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, passed in as the EXPLICIT LogGroup construct
|
|
(from ``make_function_log_group``) so CFN orders the filter after the
|
|
group exists -- a by-name import here failed the first fresh-account
|
|
deploy. 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=log_group,
|
|
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))
|