"""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: in an account where ``/aws/lambda/`` already exists out-of-band (mgmt), deploying this CREATE would collide -- acceptable because mgmt is frozen post-migration and never redeployed from main. """ 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}Alarm", alarm names f"{name_prefix}-", 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))