"""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_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, ) from constructs import Construct class PoIngestStack(Stack): def __init__(self, scope: Construct, construct_id: str, **kwargs): super().__init__(scope, construct_id, **kwargs) # --- S3 bucket for raw emails --- email_bucket = s3.Bucket( self, "EmailBucket", bucket_name=f"po-ingest-emails-{self.account}", removal_policy=RemovalPolicy.RETAIN, lifecycle_rules=[ s3.LifecycleRule(expiration=Duration.days(90)), ], ) # --- 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, ) # --- 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, ) # --- 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_retention=logs.RetentionDays.TWO_MONTHS, 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 (CMK alias/seahaven-alarm-topics lives on the # topic). Any errored invocation in a 5-min window pages. alarm_topic = sns.Topic.from_topic_arn( self, "SiteAlertsTopic", f"arn:aws:sns:{self.region}:{self.account}:site-alerts", ) 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)) # 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 = 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_retention=logs.RetentionDays.TWO_MONTHS, environment={ "PO_TABLE": "purchase-orders", }, ) po_table.grant_read_data(web_ui) # 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 = 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_retention=logs.RetentionDays.TWO_MONTHS, 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, ) ) 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)