procurement-ingest/cdk/stack.py
Adam Moussa ec416079f5 Add verified-sites pipeline via DynamoDB Streams
Enable DynamoDB Streams on purchase-orders table and add a site-extractor
Lambda that extracts Amazon facility codes and addresses from PO ship-to
data, upserting them into a new verified-sites table. Includes a backfill
script for existing POs and upgrades existing Lambdas to arm64 + 60-day
log retention.
2026-04-30 14:26:53 -04:00

181 lines
6.4 KiB
Python

"""CDK stack for the Coupa PO email ingestion pipeline."""
import aws_cdk as cdk
from aws_cdk import (
Duration,
RemovalPolicy,
Stack,
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_iam as iam,
)
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_AND_OLD_IMAGES,
)
# --- 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",
)
# --- 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/email_processor/package"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
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)
# 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/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)
# Function URL for direct access
web_url = web_ui.add_function_url(
auth_type=lambda_.FunctionUrlAuthType.NONE,
)
cdk.CfnOutput(self, "WebUIUrl", value=web_url.url, description="PO Dashboard URL")
# --- 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,
)
verified_sites_table.add_global_secondary_index(
index_name="by-state",
partition_key=dynamodb.Attribute(
name="state",
type=dynamodb.AttributeType.STRING,
),
projection_type=dynamodb.ProjectionType.ALL,
)
# --- 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/site_extractor"),
timeout=Duration.seconds(60),
memory_size=256,
log_retention=logs.RetentionDays.TWO_MONTHS,
environment={
"VERIFIED_SITES_TABLE": verified_sites_table.table_name,
},
)
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",
)