"""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/`` 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}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))