mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 08:23:14 +00:00
Some checks are pending
Deploy / deploy (push) Waiting to run
* feat: Python derived-field classifier with shadow telemetry for PO ingest Port the site_code/trade/fiscal_year rules from EXTRACTION_PROMPT into a pure, total derived_fields module applied in the shared enrich_parsed() post-stage. Python fills gaps on both parse paths (the template path has no LLM values, closing the derived-field gap opened by the PR #105 two-PR split) and never overwrites a non-null LLM value; on ai_fallback a DerivedFieldAgreement EMF record per field shadows Python against the LLM during the bake. Rules hardened against a full-corpus backtest (3,422 real emails vs the LLM-written baseline): site_code 99.4% with zero Python-wrong cases, fiscal_year 100%, trade 96.8% ex-deliberate. Also: quantity/price now declared numeric in the prompt, and derived_fields.py added to the po_stack bundling copy (deploy-time ImportError otherwise). * fix: security-review hardening — EMF value length clamp, aggregate trade CPU budget sh-security-review (4 detectors + proof-or-kill verifier): PASS, 0 confirmed critical/high. Fixes the one confirmed low (unbounded LLM-value str() into the DerivedFieldAgreement EMF log line, clamped to 64 chars) and adds the verifier-recommended defense-in-depth aggregate character budget across line items in derive_trade (per-item caps alone allowed ~10s full-core on a pathological direct-call input; unreachable through the deployed handler but cheap to bound). Bundling cp list now carries a warning comment (GPT-4.1 cross-review FIX).
701 lines
32 KiB
Python
701 lines
32 KiB
Python
"""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_iam as iam,
|
|
aws_kms as kms,
|
|
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,
|
|
aws_ssm as ssm,
|
|
)
|
|
from constructs import Construct
|
|
|
|
# Operations these tables actually issue (PutItem/UpdateItem/DeleteItem writes,
|
|
# GetItem/Query/BatchGetItem reads). DynamoDB emits ThrottledRequests/SystemErrors
|
|
# keyed by TableName + Operation only, so the CDK *_for_operations helpers (which
|
|
# render a SUM MathExpression across these per-operation metrics) are the correct,
|
|
# non-deprecated way to roll a table up to a single alarmable series.
|
|
_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,
|
|
]
|
|
|
|
|
|
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))
|
|
|
|
|
|
# CloudWatch namespace for the log-derived sender-authentication metrics.
|
|
_SENDER_AUTH_METRIC_NAMESPACE = "Seahaven/ProcurementIngest"
|
|
|
|
|
|
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/<fn>`` 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,
|
|
)
|
|
|
|
# Fire on a *sustained* reject condition rather than a volume spike. The
|
|
# earlier Sum>=3-over-15-min threshold had a blind spot that is exactly the
|
|
# failure this alarm exists to catch: a low-traffic pipeline in total
|
|
# drift outage (allowlist wrong / signing-domain changed) may only produce
|
|
# a trickle of rejects -- one every few minutes -- that never sums to 3 in
|
|
# any window, so the outage never pages. Instead: >=1 reject per 5-min
|
|
# period, alarming when 2 of the last 6 periods breach (evaluation_periods=6
|
|
# / datapoints_to_alarm=2). Six periods (30 min) with only 2 required
|
|
# datapoints closes the sparse-outage residual: even rejections >10-15 min
|
|
# apart can still place two breaching datapoints in a single 30-min
|
|
# evaluation window. A single stray spoof probe (one lone period) is
|
|
# tolerated and self-clears, but a sustained reject condition trips even
|
|
# at very low arrival rates. default_value=0 on the metric filter keeps the
|
|
# series continuous so NOT_BREACHING only applies before the first datapoint
|
|
# ever arrives.
|
|
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))
|
|
|
|
|
|
class PoIngestStack(Stack):
|
|
def __init__(self, scope: Construct, construct_id: str, **kwargs):
|
|
super().__init__(scope, construct_id, **kwargs)
|
|
|
|
# --- Shared alarm SNS topic (site-alerts) ---
|
|
# Imported once near the top so every alarm in this stack reuses the same
|
|
# Topic construct instance (avoids duplicate logical IDs). ALARM-only
|
|
# SnsAction; no OK action, per the CloudWatch-alarm preference. The
|
|
# topic's CMK (alias/seahaven-alarm-topics) lives on the topic itself.
|
|
alarm_topic = sns.Topic.from_topic_arn(
|
|
self,
|
|
"SiteAlertsTopic",
|
|
f"arn:aws:sns:{self.region}:{self.account}:site-alerts",
|
|
)
|
|
|
|
# --- S3 bucket for raw emails ---
|
|
email_bucket = s3.Bucket(
|
|
self,
|
|
"EmailBucket",
|
|
bucket_name=f"po-ingest-emails-{self.account}",
|
|
block_public_access=s3.BlockPublicAccess.BLOCK_ALL,
|
|
removal_policy=RemovalPolicy.RETAIN,
|
|
lifecycle_rules=[
|
|
s3.LifecycleRule(expiration=Duration.days(90)),
|
|
],
|
|
)
|
|
|
|
# --- Shared customer-managed CMK for sensitive DynamoDB tables ---
|
|
# Owned by the account-baseline app (alias/seahaven-dynamodb, INFRA-95 /
|
|
# M-3); ARN published to SSM. The purchase-orders table was migrated to
|
|
# SSE-KMS out-of-band, so declaring encryption_key here reconciles the
|
|
# drift and — via grant_read_write_data below — propagates the required
|
|
# kms:Decrypt/GenerateDataKey/DescribeKey to the consumer roles.
|
|
dynamodb_cmk = kms.Key.from_key_arn(
|
|
self,
|
|
"DynamoDbCmk",
|
|
ssm.StringParameter.value_for_string_parameter(
|
|
self, "/seahaven/dynamodb/cmk-arn"
|
|
),
|
|
)
|
|
|
|
# --- 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,
|
|
encryption=dynamodb.TableEncryption.CUSTOMER_MANAGED,
|
|
encryption_key=dynamodb_cmk,
|
|
)
|
|
|
|
# --- Anthropic API key secret removed (Bedrock migration) ---
|
|
# PO parsing stays fully AI but moved from the Anthropic API to the
|
|
# Bedrock inference profile us.anthropic.claude-haiku-4-5-20251001-v1:0,
|
|
# so no provider API key is needed. The old secret
|
|
# "po-ingest/anthropic-api-key" had RemovalPolicy.RETAIN, so it is
|
|
# ORPHANED (not deleted) by this change: delete it manually post-deploy
|
|
# and revoke the stored key at Anthropic.
|
|
|
|
# --- 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 && "
|
|
# NOTE: every module handler.py imports as a sibling
|
|
# MUST be listed here or the deploy ships a Lambda that
|
|
# ImportErrors at runtime (bit us for template_parser
|
|
# in PR #105 and nearly for derived_fields in PR #2).
|
|
"cp handler.py ses_auth.py template_parser.py "
|
|
"derived_fields.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",
|
|
"BEDROCK_MODEL_ID": "us.anthropic.claude-haiku-4-5-20251001-v1:0",
|
|
# Fail-closed sender auth (INFRA-107): the handler only
|
|
# accepts mail whose SES-stamped Authentication-Results
|
|
# header carries dkim=pass for one of these domains.
|
|
# Observed on live traffic 2026-07-15: Coupa PO mail passes
|
|
# DKIM for amazon.coupahost.com (and amazonses.com, which is
|
|
# deliberately NOT allowlisted — every SES customer's mail
|
|
# passes that). Unset/empty ⇒ the handler rejects all mail.
|
|
"ALLOWED_DKIM_DOMAINS": "amazon.coupahost.com",
|
|
},
|
|
)
|
|
|
|
# Grant permissions
|
|
email_bucket.grant_read(email_processor)
|
|
po_table.grant_read_write_data(email_processor)
|
|
|
|
# --- Bedrock InvokeModel grant ---
|
|
# The us.* inference profile can route cross-region, so the grant MUST
|
|
# cover both the inference-profile ARN AND the per-region foundation-model
|
|
# ARNs (empty account field) for every region the profile can reach
|
|
# (us-east-1/us-east-2/us-west-2). A profile-only grant AccessDenies at
|
|
# runtime whenever the profile routes to a region whose foundation-model
|
|
# ARN is not allowed.
|
|
email_processor.add_to_role_policy(
|
|
iam.PolicyStatement(
|
|
actions=[
|
|
"bedrock:InvokeModel",
|
|
"bedrock:InvokeModelWithResponseStream",
|
|
],
|
|
resources=[
|
|
"arn:aws:bedrock:us-east-1:328440206208:inference-profile/us.anthropic.claude-haiku-4-5-20251001-v1:0",
|
|
"arn:aws:bedrock:us-east-1::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",
|
|
],
|
|
)
|
|
)
|
|
|
|
# --- Errors alarm (INFRA-41 / audit H-8) ---
|
|
# ALARM-only (no OK action, per the CloudWatch-alarm preference) to the
|
|
# shared site-alerts topic. Any errored invocation in a 5-min window pages.
|
|
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))
|
|
|
|
# --- Sender-auth rejection alarm (INFRA-107) ---
|
|
# A rejected email (bad/unaligned DKIM verdict) returns normally, so it
|
|
# produces NO Lambda error, NO DLQ message and NO retry -- only a
|
|
# `sender_auth_rejected` warning log. Without this metric filter + alarm a
|
|
# domain drift (Coupa rotates its signing subdomain, SES changes its
|
|
# Authentication-Results format, the allowlist is wrong) would silently
|
|
# discard 100% of legitimate PO mail while every other alarm stays green.
|
|
# A CloudWatch Logs metric filter turns those warnings into a metric so a
|
|
# false-reject storm pages instead of vanishing. default_value=0 keeps the
|
|
# series populated (alarm stays OK, never INSUFFICIENT_DATA) between events.
|
|
_add_sender_auth_rejected_alarm(
|
|
self, "EmailProcessor", "po-email-processor", alarm_topic
|
|
)
|
|
|
|
# --- Throttles alarm: po-email-processor ---
|
|
# Any throttled invocation (concurrency cap hit) in a 5-min window pages.
|
|
# ALARM-only to site-alerts; no OK action; NOT_BREACHING when no data.
|
|
email_processor.metric_throttles(
|
|
period=Duration.minutes(5),
|
|
statistic="Sum",
|
|
).create_alarm(
|
|
self,
|
|
"EmailProcessorThrottlesAlarm",
|
|
alarm_name="po-email-processor-throttles",
|
|
alarm_description="po-email-processor invocation 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(alarm_topic))
|
|
|
|
# --- DLQ messages-present alarm ---
|
|
# Pages when any message lands in the EmailProcessorDlq: a message here
|
|
# means a PO email was permanently dropped after Lambda exhausted its
|
|
# async retries. Maximum over a single 5-min window > 0 fires; missing
|
|
# data (no messages metric emitted) is not breaching. Reuses the shared
|
|
# site-alerts topic, ALARM-only, like the errors alarm above. The metric
|
|
# helper derives the QueueName dimension from the queue construct, so the
|
|
# alarm tracks the CDK-generated queue name without hardcoding it.
|
|
email_processor_dlq.metric_approximate_number_of_messages_visible(
|
|
period=Duration.minutes(5),
|
|
statistic="Maximum",
|
|
).create_alarm(
|
|
self,
|
|
"EmailProcessorDlqMessagesAlarm",
|
|
alarm_name="po-email-processor-dlq-messages",
|
|
alarm_description="po-email-processor DLQ has messages (dropped PO emails)",
|
|
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))
|
|
|
|
# --- Duration alarm: po-email-processor (orphan adoption) ---
|
|
# Adopts the orphaned CLI alarm Lambda-Duration-po-email-processor under
|
|
# the repo's <fn>-duration naming (NEW logical name → no deploy collision;
|
|
# delete the orphan post-deploy). p99 / 45000 ms
|
|
# (75% of the 60s timeout) / eval 3 of 3 — tighter than the orphan's
|
|
# Maximum>=48000 / 1-of-1.
|
|
email_processor.metric_duration(
|
|
period=Duration.minutes(5),
|
|
statistic="p99",
|
|
).create_alarm(
|
|
self,
|
|
"EmailProcessorDurationAlarm",
|
|
alarm_name="po-email-processor-duration",
|
|
alarm_description="po-email-processor p99 duration approaching the 60s timeout",
|
|
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(alarm_topic))
|
|
|
|
# --- Template fallback-rate alarm: po-email-processor ---
|
|
# The processor tries a deterministic template parse first and only calls
|
|
# the Bedrock AI extractor on a miss/invalid. A sustained rise in the
|
|
# ai_fallback share signals Coupa template drift (coverage collapse).
|
|
# EMF metric Seahaven/PoIngest/ParseOutcome, dimensioned by ParseMethod
|
|
# (template|ai_fallback).
|
|
#
|
|
# RETUNED for PO volume (~57 emails/day ≈ 14.25 per 6h period) -- the WO
|
|
# alarm's 15-min period / >=10-sample floor assume ~760/day and would be
|
|
# structurally DEAD here (a 15-min period holds ~0.6 PO emails, so the
|
|
# floor is never met and the IF always takes the 0 branch):
|
|
# * period 6h: a stable ~14-email denominator per datapoint.
|
|
# * volume floor >=8: at the floor, one fallback email = 12.5% < 20%,
|
|
# so a single email can NEVER breach a datapoint; a breach needs >=2
|
|
# fallbacks in one 6h window (2/8 = 25%) or >=3 at typical volume
|
|
# (3/14 ≈ 21%). Sparse overnight/weekend windows (<8 emails) take
|
|
# the 0 branch -- non-breaching by design (accepted trade: a Friday-
|
|
# evening drift may not page until weekend volume accrues).
|
|
# * threshold >20%: expected baseline fallback ≈1% (comments 0.55% +
|
|
# multi-line 0.18% + non-USD 0) -- far below the threshold.
|
|
# * 2 of 4 datapoints (24h span): isolated noise self-clears, while
|
|
# total template drift (100% fallback) pages within ~12h.
|
|
# Post-#102 rule: NO element-wise MAX(timeseries, scalar) in alarm math;
|
|
# the IF volume floor guarantees the non-zero denominator. Any change to
|
|
# this expression must be gated by `npx cdk synth po-ingest`.
|
|
fb_metric = cloudwatch.Metric(
|
|
namespace="Seahaven/PoIngest",
|
|
metric_name="ParseOutcome",
|
|
dimensions_map={"ParseMethod": "ai_fallback"},
|
|
statistic="Sum",
|
|
period=Duration.hours(6),
|
|
)
|
|
tmpl_metric = cloudwatch.Metric(
|
|
namespace="Seahaven/PoIngest",
|
|
metric_name="ParseOutcome",
|
|
dimensions_map={"ParseMethod": "template"},
|
|
statistic="Sum",
|
|
period=Duration.hours(6),
|
|
)
|
|
fallback_rate = cloudwatch.MathExpression(
|
|
expression=(
|
|
"IF((FILL(fb,0)+FILL(tmpl,0))>=8, "
|
|
"100*FILL(fb,0)/(FILL(fb,0)+FILL(tmpl,0)), 0)"
|
|
),
|
|
using_metrics={"fb": fb_metric, "tmpl": tmpl_metric},
|
|
period=Duration.hours(6),
|
|
label="TemplateFallbackRatePct",
|
|
)
|
|
fallback_rate.create_alarm(
|
|
self,
|
|
"EmailProcessorTemplateFallbackRateAlarm",
|
|
alarm_name="po-email-processor-template-fallback-rate",
|
|
alarm_description=(
|
|
"po-email-processor deterministic-template coverage collapse: "
|
|
">20% of parses fell back to the Bedrock AI extractor"
|
|
),
|
|
threshold=20,
|
|
evaluation_periods=4,
|
|
datapoints_to_alarm=2,
|
|
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 auth token secret ---
|
|
# Shared secret for the web UI auth gate, stored in Secrets Manager and
|
|
# resolved at runtime so the token never appears in CloudFormation templates
|
|
# or Lambda environment variables. Create this secret before deploying
|
|
# either stack; both PO and WO stacks reference it by name.
|
|
web_ui_auth_secret = secretsmanager.Secret.from_secret_name_v2(
|
|
self,
|
|
"WebUiAuthToken",
|
|
"procurement-ingest/web-ui-auth-token",
|
|
)
|
|
|
|
# --- 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",
|
|
# Defense-in-depth shared secret for the web UI handler. The
|
|
# handler fails closed if this ARN is unset or the secret is
|
|
# missing, so any future invocation path cannot re-expose the
|
|
# PO DB unauthenticated. The secret value is fetched at runtime
|
|
# from Secrets Manager (not embedded in env vars or template).
|
|
"WEB_UI_AUTH_TOKEN_SECRET_ARN": web_ui_auth_secret.secret_arn,
|
|
},
|
|
)
|
|
|
|
po_table.grant_read_data(web_ui)
|
|
web_ui_auth_secret.grant_read(web_ui)
|
|
|
|
# --- Throttles alarm: po-web-ui ---
|
|
web_ui.metric_throttles(
|
|
period=Duration.minutes(5),
|
|
statistic="Sum",
|
|
).create_alarm(
|
|
self,
|
|
"WebUiThrottlesAlarm",
|
|
alarm_name="po-web-ui-throttles",
|
|
alarm_description="po-web-ui invocation 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(alarm_topic))
|
|
|
|
# --- Duration alarm: po-web-ui ---
|
|
# Net-new (no orphan exists for this function).
|
|
# p99 / 45000 ms (75% of the 60s timeout) / eval 3, datapoints 2.
|
|
web_ui.metric_duration(
|
|
period=Duration.minutes(5),
|
|
statistic="p99",
|
|
).create_alarm(
|
|
self,
|
|
"WebUiDurationAlarm",
|
|
alarm_name="po-web-ui-duration",
|
|
alarm_description="po-web-ui p99 duration approaching the 60s timeout",
|
|
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(alarm_topic))
|
|
|
|
# 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,
|
|
)
|
|
)
|
|
|
|
# --- Errors alarm: po-ingest-site-extractor ---
|
|
# Stream-consumer errors retry per the event-source config, but a
|
|
# persistent failure stalls the verified-sites pipeline. ALARM-only to
|
|
# site-alerts; no OK action; NOT_BREACHING when no data.
|
|
site_extractor.metric_errors(
|
|
period=Duration.minutes(5),
|
|
statistic="Sum",
|
|
).create_alarm(
|
|
self,
|
|
"SiteExtractorErrorsAlarm",
|
|
alarm_name="po-ingest-site-extractor-errors",
|
|
alarm_description="po-ingest-site-extractor 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))
|
|
|
|
# --- Throttles alarm: po-ingest-site-extractor ---
|
|
site_extractor.metric_throttles(
|
|
period=Duration.minutes(5),
|
|
statistic="Sum",
|
|
).create_alarm(
|
|
self,
|
|
"SiteExtractorThrottlesAlarm",
|
|
alarm_name="po-ingest-site-extractor-throttles",
|
|
alarm_description="po-ingest-site-extractor invocation 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(alarm_topic))
|
|
|
|
# --- Duration alarm: po-ingest-site-extractor ---
|
|
# Net-new (no orphan exists for this function).
|
|
# p99 / 45000 ms (75% of the 60s timeout) / eval 3, datapoints 2.
|
|
site_extractor.metric_duration(
|
|
period=Duration.minutes(5),
|
|
statistic="p99",
|
|
).create_alarm(
|
|
self,
|
|
"SiteExtractorDurationAlarm",
|
|
alarm_name="po-ingest-site-extractor-duration",
|
|
alarm_description="po-ingest-site-extractor p99 duration approaching the 60s timeout",
|
|
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(alarm_topic))
|
|
|
|
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)
|
|
|
|
# --- DynamoDB throttle + system-error alarms ---
|
|
# ThrottledRequests / SystemErrors emit at TableName + Operation only
|
|
# (verified against live CloudWatch: no TableName-only rollup exists, and
|
|
# metric_throttled_requests is deprecated/invalid in aws-cdk-lib 2.259.0).
|
|
# Each table currently has zero throttle/error datapoints, so the series
|
|
# only materialise on first occurrence — NOT_BREACHING keeps them OK until
|
|
# then.
|
|
_add_ddb_alarms(
|
|
self, "PurchaseOrdersTable", po_table, "purchase-orders", alarm_topic
|
|
)
|
|
_add_ddb_alarms(
|
|
self,
|
|
"VerifiedSitesTable",
|
|
verified_sites_table,
|
|
"verified-sites",
|
|
alarm_topic,
|
|
)
|
|
_add_ddb_alarms(
|
|
self,
|
|
"PendingSiteReviewTable",
|
|
pending_review_table,
|
|
"pending-site-review",
|
|
alarm_topic,
|
|
)
|