"""CDK stack for the work order email ingestion pipeline.""" import aws_cdk as cdk import common from aws_cdk import ( Duration, RemovalPolicy, Stack, ) from aws_cdk import ( aws_cloudwatch as cloudwatch, ) from aws_cdk import ( aws_cloudwatch_actions as cw_actions, ) from aws_cdk import ( aws_dynamodb as dynamodb, ) from aws_cdk import ( aws_iam as iam, ) from aws_cdk import ( aws_kms as kms, ) from aws_cdk import ( aws_lambda as lambda_, ) from aws_cdk import ( aws_lambda_event_sources as lambda_event_sources, ) from aws_cdk import ( aws_s3 as s3, ) from aws_cdk import ( aws_s3_notifications as s3n, ) from aws_cdk import ( aws_secretsmanager as secretsmanager, ) from aws_cdk import ( aws_ses as ses, ) from aws_cdk import ( aws_ses_actions as ses_actions, ) from aws_cdk import ( aws_sns as sns, ) from aws_cdk import ( aws_sqs as sqs, ) from constructs import Construct class WorkorderIngestStack(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 = common.make_email_bucket( self, "EmailBucket", "workorder-ingest-emails" ) # --- DynamoDB tables --- work_orders_table = dynamodb.Table( self, "WorkOrdersTable", table_name="WorkOrders", partition_key=dynamodb.Attribute( name="work_order_id", type=dynamodb.AttributeType.STRING, ), billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, # NEW_AND_OLD_IMAGES: the SHOC emitter needs OLD.wo_status to # classify the cancelled transition (docs/shoc-webhook-plan.md # Phase 2). In-place CFN update -- no table replacement. stream=dynamodb.StreamViewType.NEW_AND_OLD_IMAGES, removal_policy=RemovalPolicy.RETAIN, ) # site-code-index and status-index GSIs removed 2026-06-03 (audit M-20): # 0 reads in 30d against ~50k WCU each of write amplification. Re-add if # a site-code or status query path ships. comments_table = dynamodb.Table( self, "CommentsTable", table_name="WorkOrderComments", partition_key=dynamodb.Attribute( name="work_order_id", type=dynamodb.AttributeType.STRING, ), sort_key=dynamodb.Attribute( name="comment_id", type=dynamodb.AttributeType.STRING, ), billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, # Streamed for the SHOC emitter (docs/shoc-webhook-plan.md Phase 2). stream=dynamodb.StreamViewType.NEW_AND_OLD_IMAGES, removal_policy=RemovalPolicy.RETAIN, ) # --- Anthropic API key secret removed (Bedrock migration) --- # Parsing 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 "workorder-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 workorder-email-processor-dlq (removed post-deploy). email_processor_dlq = common.make_processor_dlq(self, "EmailProcessorDlq") # --- Lambda function --- email_processor_log_group = common.make_function_log_group( self, "EmailProcessor", "workorder-email-processor" ) email_processor = lambda_.Function( self, "EmailProcessor", function_name="workorder-email-processor", runtime=lambda_.Runtime.PYTHON_3_12, architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( "../lambdas", exclude=["**/__pycache__/**", "**/tests/**", "**/package/**"], bundling=cdk.BundlingOptions( image=lambda_.Runtime.PYTHON_3_12.bundling_image, command=[ "bash", "-c", # pip step removed in Phase 7: requirements.txt is now empty # (boto3 comes from the Lambda runtime), so nothing is installed # and the manylinux pin has nothing to pin. cp-only is safe. # shared/*.py ships the four modules extracted to # lambdas/shared/ (Phase 3): ses_auth, web_ui_auth, # email_parsing, emf. Flat cp keeps the bare-name # imports (e.g. `from ses_auth import ...`) resolving # unchanged in /asset-output. "cp wo/email_processor/*.py /asset-output/ && " "cp shared/*.py /asset-output/", ], ), ), timeout=Duration.seconds(60), memory_size=256, log_group=email_processor_log_group, dead_letter_queue=email_processor_dlq, environment={ "WORK_ORDERS_TABLE": work_orders_table.table_name, "COMMENTS_TABLE": comments_table.table_name, "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. APM # mail arrives via the apm@ Google Groups forward, which # re-signs as seahaven.com (observed on live traffic # 2026-07-15: "dkim=pass header.i=@seahaven.com"; the # original hxgnsmartcloud.com signature does not survive # the forward). Unset/empty ⇒ the handler rejects all mail. "ALLOWED_DKIM_DOMAINS": "seahaven.com", }, ) # Grant permissions email_bucket.grant_read(email_processor) work_orders_table.grant_read_write_data(email_processor) comments_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(common.make_bedrock_invoke_statement(self)) # NOTE: The pre-emptive grant_encrypt_decrypt on the shared DynamoDB CMK # (alias/seahaven-dynamodb) was removed (security sweep 2026-06-17). The # WorkOrders/WorkOrderComments tables are NOT SSE-KMS encrypted with that # CMK, so the grant was unused for these tables yet handed # wo-email-processor kms:Decrypt on the CMK that also protects the # purchase-orders table (cross-stack decrypt reach). Re-add this grant only # as part of the actual CMK migration of these tables (INFRA-6), at which # point grant_read_write_data on the (then encrypted) tables would propagate # the needed key permissions automatically. # --- Standard per-Lambda alarms: workorder-email-processor --- # errors (INFRA-41 / audit H-8), throttles, DLQ-visible-messages (dropped # emails), and a p95 duration alarm (orphan adoption of the CLI # Lambda-Duration-workorder-email-processor under -duration naming, # 45000 ms = 75% of the 60s timeout, eval 3 / dp 2). p95 (NOT p99) is the # WO-specific duration statistic. All ALARM-only to site-alerts. common.add_standard_lambda_alarms( self, "EmailProcessor", email_processor, "workorder-email-processor", alarm_topic, duration_statistic="p95", errors=True, dlq=email_processor_dlq, descriptions={ "errors": "workorder-email-processor async invocation errors", "throttles": "workorder-email-processor invocation throttles", "dlq": "workorder-email-processor DLQ has visible messages (dropped emails)", "duration": "workorder-email-processor p95 duration approaching the 60s timeout", }, ) # --- 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. The WO allowlist trusts dkim=pass # for seahaven.com on the assumption the apm@ forward re-signs there; if # that assumption is wrong (e.g. a Gmail auto-forward re-signs under a # different domain), 100% of legitimate work-order mail is silently # dropped. This metric filter + alarm turns those warnings into a paging # signal so a false-reject storm surfaces instead of a silent outage. common.add_sender_auth_rejected_alarm( self, "EmailProcessor", "workorder-email-processor", alarm_topic, email_processor_log_group, ) # --- Template fallback-rate alarm: workorder-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 Hexagon template drift (coverage collapse). # EMF metric Seahaven/WorkorderIngest/ParseOutcome, dimensioned by # ParseMethod (template|ai_fallback). 15-min periods (deliberate deviation # from the 5-min house style) accumulate a stable denominator at the low # ~760/day volume; FILL(0) + a >=10-sample volume floor prevent # low-volume false pages and INSUFFICIENT_DATA. ALARM-only SnsAction to # site-alerts, no OK action, NOT_BREACHING -- matching the stack idiom. # rejected_included=True: the WO expression folds ai_fallback_rejected # (rej) into BOTH numerator and denominator -- a drift outage whose AI # output also fails the gate must still count as fallback, otherwise it # would LOWER the observed rate while silently dropping mail. (Contrast # PO, which excludes rej to avoid a pre-call double-count.) common.make_fallback_rate_alarm( self, "EmailProcessorTemplateFallbackRateAlarm", namespace="Seahaven/WorkorderIngest", alarm_topic=alarm_topic, alarm_name="workorder-email-processor-template-fallback-rate", alarm_description=( "workorder-email-processor deterministic-template coverage " "collapse: >15% of parses fell back to the Bedrock AI extractor" ), rejected_included=True, period=Duration.minutes(15), threshold=15, floor=10, evaluation_periods=3, datapoints_to_alarm=2, ) # --- AI-fallback rejected alarm: workorder-email-processor --- # A parse rejected by the validate_ai_fallback gate is dropped without # error/retry/DLQ (fail closed), so like sender-auth rejections it # needs its own pager or a sustained rejection condition (prompt- # injection probing, or template drift whose AI output fails the gate) # stays silent. Same sparse-arrival idiom as the sender-auth-rejected # alarm: >=1 rejection per 5-min period, 2 of the last 6 periods (30 # min), so a lone probe self-clears but a burst pages within ~10 min. # Coverage residual (matching the sender-auth-rejected sibling and # knowingly accepted): rejections spaced >~25-30 min apart never place # two breaching datapoints in one 30-min window, and the fallback-rate # alarm dilutes them below 15% against normal template volume, so a # *very* sparse silent-drop trickle is not paged by either alarm. # EMF emits no datapoint in quiet periods (no metric-filter # default_value here); NOT_BREACHING treats those gaps as OK. cloudwatch.Metric( namespace="Seahaven/WorkorderIngest", metric_name="ParseOutcome", dimensions_map={"ParseMethod": "ai_fallback_rejected"}, statistic="Sum", period=Duration.minutes(5), ).create_alarm( self, "EmailProcessorAiFallbackRejectedAlarm", alarm_name="workorder-email-processor-ai-fallback-rejected", alarm_description=( "workorder-email-processor is rejecting Bedrock AI-fallback " "output at the validation gate (possible prompt-injection " "probing or template 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)) # S3 event notification -> Lambda email_bucket.add_event_notification( s3.EventType.OBJECT_CREATED, s3n.LambdaDestination(email_processor), s3.NotificationKeyFilter(prefix="inbound/"), ) # --- SES Receipt Rule --- rule_set = ses.ReceiptRuleSet.from_receipt_rule_set_name( self, "ExistingRuleSet", "INBOUND_MAIL", ) rule_set.add_rule( "WorkorderEmailRule", recipients=["apm@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_log_group = common.make_function_log_group( self, "WebUI", "workorder-web-ui" ) web_ui = lambda_.Function( self, "WebUI", function_name="workorder-web-ui", runtime=lambda_.Runtime.PYTHON_3_12, architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( "../lambdas", exclude=["**/__pycache__/**"], bundling=cdk.BundlingOptions( image=lambda_.Runtime.PYTHON_3_12.bundling_image, command=[ "bash", "-c", # web_ui_auth.py is shared (lambdas/shared) and must land # FLAT beside handler.py so # `from web_ui_auth import is_authenticated` resolves at # runtime. Only web_ui_auth is copied from shared/. WO # baseline keeps __init__.py and requirements.txt, so the # whole web_ui dir is copied; deployed file list becomes # {__init__, handler, requirements.txt, web_ui_auth}. # NOTE: the top-level `exclude=` on from_asset only # filters the asset-hash fingerprint, NOT the directory # Docker bundling actually mounts, so a local # __pycache__ on disk at synth time WOULD otherwise leak # into the bundled zip -- strip it explicitly post-cp # instead of relying on exclude. "cp -r wo/web_ui/. /asset-output/ && " "cp shared/web_ui_auth.py /asset-output/ && " "rm -rf /asset-output/__pycache__", ], ), ), timeout=Duration.seconds(15), memory_size=128, log_group=web_ui_log_group, environment={ "WORK_ORDERS_TABLE": work_orders_table.table_name, "COMMENTS_TABLE": comments_table.table_name, # 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 # WO 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, }, ) work_orders_table.grant_read_data(web_ui) comments_table.grant_read_data(web_ui) web_ui_auth_secret.grant_read(web_ui) # --- 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.261.0). # Each table currently has zero throttle/error datapoints, so the series # only materialise on first occurrence — NOT_BREACHING keeps them OK until # then. common.add_ddb_alarms( self, "WorkOrdersTable", work_orders_table, "WorkOrders", alarm_topic ) common.add_ddb_alarms( self, "WorkOrderComments", comments_table, "WorkOrderComments", alarm_topic ) # --- Function ARN + consumed-table-name outputs (Phase 4, additive) --- cdk.CfnOutput( self, "EmailProcessorFunctionArn", value=email_processor.function_arn, description="ARN of the workorder-email-processor Lambda", ) cdk.CfnOutput( self, "WebUiFunctionArn", value=web_ui.function_arn, description="ARN of the workorder-web-ui Lambda", ) cdk.CfnOutput( self, "WorkOrdersTableName", value=work_orders_table.table_name, description="WorkOrders DynamoDB table", ) cdk.CfnOutput( self, "WorkOrderCommentsTableName", value=comments_table.table_name, description="WorkOrderComments DynamoDB table", ) # 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. # ===================================================================== # --- SHOC webhook emitter (docs/shoc-webhook-plan.md) --- # Realtime work-order feed to the SHOC backend: DynamoDB Streams on the # two WO tables -> workorder-shoc-emitter -> HMAC-signed HTTPS POST. # Contract: docs/shoc-webhook-contract.md (Rev 2026-07-23). Built in # _add_shoc_webhook_emitter (module helper, common.py plain-helper # style) to keep __init__ under the PLR0915 statement ceiling; the # stack stays the construct scope, so extraction does not move any # logical ID. # ===================================================================== _add_shoc_webhook_emitter(self, work_orders_table, comments_table, alarm_topic) def _scope_rotation_invoke_permission(stack, rotator_fn, secret): """Add SourceAccount/SourceArn to the generated rotation invoke permission. ``add_rotation_schedule`` emits an ``AWS::Lambda::Permission`` for the ``secretsmanager.amazonaws.com`` service principal with no source conditions. Rather than add a second (additive) permission, find that generated ``CfnPermission`` and pin it to this account + secret so only this secret's Secrets Manager can invoke the rotator. """ patched = False for child in stack.node.find_all(): if ( isinstance(child, lambda_.CfnPermission) and child.principal == "secretsmanager.amazonaws.com" and stack.resolve(child.function_name) == stack.resolve(rotator_fn.function_arn) ): child.source_account = stack.account child.source_arn = secret.secret_arn patched = True if not patched: # fail loud if CDK changes the generated shape on upgrade raise RuntimeError( "rotation invoke CfnPermission not found; cannot scope source conditions" ) def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_topic): """SHOC webhook emitter (docs/shoc-webhook-plan.md Phases 1-4). Dedicated HMAC CMK + secret with cross-account SHOC read grants, the 30-day rotation Lambda, the two failure queues, the stream-driven emitter Lambda with its two (dark, enabled=False) event source mappings, and the full alarm set. `stack` is the construct scope for every child, exactly as if this code were inline in __init__. """ # SHOC consumer principal for the cross-account read grants. EXACT role # ARN only -- future shoc-backend-staging/-prod roles are each a # deliberate, individually-reviewed policy addition (no wildcard or # prefix trust). shoc_consumer_principal = iam.ArnPrincipal( "arn:aws:iam::396287094661:role/shoc-backend-dev" ) # --- Dedicated CMK for the HMAC secret --- # Dedicated key, NOT alias/seahaven-dynamodb: reusing the DynamoDB CMK # would hand the SHOC cross-account grant decrypt reach over the PO # table's encryption key -- the dedicated key scopes the grant to # exactly this secret (plan Phase 1). shoc_webhook_key = kms.Key( stack, "ShocWebhookHmacKey", alias="workorder-ingest-shoc-webhook-kms", description=( "Dedicated CMK for the workorder-ingest/shoc-webhook-hmac " "secret (cross-account readable by the SHOC backend)" ), enable_key_rotation=True, # GOTCHA: DESTROY is deliberate -- do not "harden" this to RETAIN. # The key protects only machine-generated HMAC material that is # fully regenerable by one rotation, and DESTROY avoids the # fixed-name RETAIN-orphan deadlock on the alias (mirrors the # secret's rationale below). removal_policy=RemovalPolicy.DESTROY, pending_window=Duration.days(7), ) # Key-policy half of the cross-account read grant (the secret resource # policy below is the other half; either one alone fails silently at # the receiver). resources=["*"] is key-scoped, not account-wide -- # KMS key policies only ever apply to this key. # kms:ViaService pins the grant to Secrets Manager decrypt paths only # (GPT-4.1 cross-review FIX): a compromised shoc-backend-dev cannot use # this key for arbitrary KMS operations outside the secret fetch. shoc_webhook_key.add_to_resource_policy( iam.PolicyStatement( actions=["kms:Decrypt"], principals=[shoc_consumer_principal], resources=["*"], conditions={ "StringEquals": { "kms:ViaService": (f"secretsmanager.{stack.region}.amazonaws.com") } }, ) ) # --- HMAC signing secret --- # Value shape (contract section 6.1): # {"keys": [{"kid": "", "secret": "<64 hex>"}, ...]}, # newest first, max 2; the producer signs with keys[0]. The # generate_secret_string below is BOOTSTRAP shape only ({"keys": []} # plus throwaway entropy the rotator ignores); the first rotation # (rotate_immediately default) populates the real keys. shoc_hmac_secret = secretsmanager.Secret( stack, "ShocWebhookHmacSecret", secret_name="workorder-ingest/shoc-webhook-hmac", encryption_key=shoc_webhook_key, description=( "HMAC signing keys for the SHOC work-order webhook " "(docs/shoc-webhook-contract.md section 6)" ), # GOTCHA: DESTROY is deliberate -- do not "harden" this to RETAIN. # The value is machine-generated HMAC material with no operator-set # content, fully regenerable by one rotation, so RETAIN buys # nothing and would expose the fixed-name RETAIN orphan deadlock # (a failed first create orphans an empty shell holding the global # name; see reference_secret_retain_orphan_deadlock). removal_policy=RemovalPolicy.DESTROY, generate_secret_string=secretsmanager.SecretStringGenerator( secret_string_template='{"keys": []}', generate_string_key="bootstrap_entropy", password_length=32, exclude_punctuation=True, ), ) cdk.Tags.of(shoc_hmac_secret).add("Purpose", "shoc-webhook-hmac") cdk.Tags.of(shoc_hmac_secret).add("ManagedBy", "procurement-ingest-cdk") # Secret-resource-policy half of the cross-account read grant (the key # policy above is the other half). DescribeSecret lets the receiver # resolve secret metadata without any broader list permission. shoc_hmac_secret.add_to_resource_policy( iam.PolicyStatement( actions=[ "secretsmanager:GetSecretValue", "secretsmanager:DescribeSecret", ], principals=[shoc_consumer_principal], resources=["*"], ) ) # --- HMAC rotation Lambda --- # 30-day schedule: generates a new key, prepends as keys[0], truncates # to 2 entries. Single-user rotation (receivers re-fetch on a <=5-min # TTL), so the standard 4-step rotation collapses to # createSecret/finishSecret. shoc_hmac_rotator_log_group = common.make_function_log_group( stack, "ShocHmacRotator", "workorder-shoc-hmac-rotator" ) shoc_hmac_rotator = lambda_.Function( stack, "ShocHmacRotator", function_name="workorder-shoc-hmac-rotator", runtime=lambda_.Runtime.PYTHON_3_12, architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( "../lambdas", exclude=["**/__pycache__/**", "**/tests/**", "**/package/**"], bundling=cdk.BundlingOptions( image=lambda_.Runtime.PYTHON_3_12.bundling_image, command=[ "bash", "-c", # Stdlib + boto3-from-runtime only; nothing installed. "cp wo/shoc_hmac_rotator/*.py /asset-output/", ], ), ), timeout=Duration.seconds(60), memory_size=128, log_group=shoc_hmac_rotator_log_group, ) # Rotation permissions, scoped to the one secret. CDK's secret_arn # token resolves to the full ARN including the -?????? suffix wildcard, # so no separate "*"-suffixed resource variant is needed. shoc_hmac_rotator.add_to_role_policy( iam.PolicyStatement( actions=[ "secretsmanager:DescribeSecret", "secretsmanager:GetSecretValue", "secretsmanager:PutSecretValue", "secretsmanager:UpdateSecretVersionStage", ], resources=[shoc_hmac_secret.secret_arn], ) ) # Explicit statement instead of grant_encrypt_decrypt so the grant can # carry kms:ViaService (cross-review FIX): the rotator only ever touches # this key through Secrets Manager put/get, never the KMS API directly. shoc_hmac_rotator.add_to_role_policy( iam.PolicyStatement( actions=[ "kms:Decrypt", "kms:Encrypt", "kms:GenerateDataKey*", "kms:ReEncrypt*", ], resources=[shoc_webhook_key.key_arn], conditions={ "StringEquals": { "kms:ViaService": (f"secretsmanager.{stack.region}.amazonaws.com") } }, ) ) shoc_hmac_secret.add_rotation_schedule( "Rotation", rotation_lambda=shoc_hmac_rotator, automatically_after=Duration.days(30), ) # Scope the Secrets-Manager-service invoke permission to THIS secret # (cross-review FIX / confused-deputy): add_rotation_schedule emits an # AWS::Lambda::Permission for secretsmanager.amazonaws.com with no # SourceAccount/SourceArn, so any account's Secrets Manager could invoke # the rotator by pointing a foreign secret's RotationLambdaARN at it. # Lambda permissions are additive (OR), so a second scoped permission # would NOT revoke the unscoped one -- patch the generated permission in # place. source_arn pins the invoker to this secret; source_account is the # belt-and-braces account bound. (Blast radius was already contained by # the rotator role being resource-scoped to this secret, but this closes # the unauthenticated invoke primitive per AWS rotation guidance.) _scope_rotation_invoke_permission(stack, shoc_hmac_rotator, shoc_hmac_secret) # --- Standard per-Lambda alarms: workorder-shoc-hmac-rotator --- # errors + throttles + p99 duration. No DLQ alarm: rotation is invoked # synchronously by Secrets Manager (dlq=None); a failed rotation # surfaces as an invocation error. common.add_standard_lambda_alarms( stack, "ShocHmacRotator", shoc_hmac_rotator, "workorder-shoc-hmac-rotator", alarm_topic, duration_statistic="p99", errors=True, dlq=None, descriptions={ "errors": "workorder-shoc-hmac-rotator invocation errors", "throttles": "workorder-shoc-hmac-rotator invocation throttles", "duration": ( "workorder-shoc-hmac-rotator p99 duration approaching the 60s timeout" ), }, ) # --- Emitter failure queues --- # Failures queue: ESM on_failure destination. It receives ESM failure # METADATA (shard/sequence pointers), not full payloads -- replay # rebuilds events from DynamoDB (contract section 8). shoc_emitter_failures_queue = sqs.Queue( stack, "ShocEmitterFailuresQueue", queue_name="workorder-shoc-emitter-failures", retention_period=Duration.days(14), enforce_ssl=True, ) # Rejected queue: full {envelope, response_status} payloads parked by # the handler on non-retryable 4xx responses (contract section 7). shoc_emitter_rejected_queue = sqs.Queue( stack, "ShocEmitterRejectedQueue", queue_name="workorder-shoc-emitter-rejected", retention_period=Duration.days(14), enforce_ssl=True, ) # --- Emitter Lambda --- shoc_emitter_log_group = common.make_function_log_group( stack, "ShocEmitter", "workorder-shoc-emitter" ) shoc_emitter = lambda_.Function( stack, "ShocEmitter", function_name="workorder-shoc-emitter", runtime=lambda_.Runtime.PYTHON_3_12, architecture=lambda_.Architecture.ARM_64, handler="handler.handler", code=lambda_.Code.from_asset( "../lambdas", exclude=["**/__pycache__/**", "**/tests/**", "**/package/**"], bundling=cdk.BundlingOptions( image=lambda_.Runtime.PYTHON_3_12.bundling_image, command=[ "bash", "-c", # Stdlib HTTP (urllib.request) + boto3-from-runtime # only; nothing installed. "cp wo/shoc_emitter/*.py /asset-output/", ], ), ), timeout=Duration.seconds(60), memory_size=256, log_group=shoc_emitter_log_group, environment={ # Non-sensitive endpoint URL (HMAC is the auth, the URL is # not). path TBD by SHOC -- confirmed in the activation PR. "SHOC_WEBHOOK_URL": ( "https://api.dev.seahaven.com/api/webhooks/work-orders" ), "HMAC_SECRET_ARN": shoc_hmac_secret.secret_arn, "REJECTED_QUEUE_URL": shoc_emitter_rejected_queue.queue_url, }, ) work_orders_table.grant_stream_read(shoc_emitter) comments_table.grant_stream_read(shoc_emitter) shoc_hmac_secret.grant_read(shoc_emitter) # Explicit statement instead of grant_decrypt so the grant carries # kms:ViaService (cross-review FIX): the emitter only decrypts this key # through Secrets Manager GetSecretValue. shoc_emitter.add_to_role_policy( iam.PolicyStatement( actions=["kms:Decrypt"], resources=[shoc_webhook_key.key_arn], conditions={ "StringEquals": { "kms:ViaService": (f"secretsmanager.{stack.region}.amazonaws.com") } }, ) ) shoc_emitter_rejected_queue.grant_send_messages(shoc_emitter) # ------------------------------------------------------------------ # ACTIVATED (enabled=True) 2026-07-30 after SHOC's receiver passed the # shared HMAC test vectors. Shipped DARK originally (enabled=False) so the # stack could deploy and be tested with zero deliveries while SHOC had no # receiver; this activation PR is the deliberate one-line flip (plan Phase # 3). LATEST start position => the feed begins now, no historical flood; # SHOC backfills history via the procurement read API, not the stream. # Ordering knobs: parallelization_factor=1, bisect_batch_on_error= # False and retry_attempts=-1 (retry until the 24h record age) are # REQUIRED for strict per-work-order in-order delivery -- a retryable # failure blocks the shard rather than skipping ahead, and # report_batch_item_failures keeps earlier in-batch successes from # being re-delivered. # ------------------------------------------------------------------ shoc_emitter.add_event_source( lambda_event_sources.DynamoEventSource( work_orders_table, starting_position=lambda_.StartingPosition.LATEST, batch_size=10, bisect_batch_on_error=False, retry_attempts=-1, max_record_age=Duration.hours(24), parallelization_factor=1, report_batch_item_failures=True, enabled=True, on_failure=lambda_event_sources.SqsDlq(shoc_emitter_failures_queue), ) ) shoc_emitter.add_event_source( lambda_event_sources.DynamoEventSource( comments_table, starting_position=lambda_.StartingPosition.LATEST, batch_size=10, bisect_batch_on_error=False, retry_attempts=-1, max_record_age=Duration.hours(24), parallelization_factor=1, report_batch_item_failures=True, enabled=True, on_failure=lambda_event_sources.SqsDlq(shoc_emitter_failures_queue), ) ) # --- Standard per-Lambda alarms: workorder-shoc-emitter --- # errors + throttles + p99 duration. No DLQ alarm here: the emitter is # a stream consumer with no async DLQ (dlq=None); its failure surfaces # are the two SQS queues alarmed bespoke below. common.add_standard_lambda_alarms( stack, "ShocEmitter", shoc_emitter, "workorder-shoc-emitter", alarm_topic, duration_statistic="p99", errors=True, dlq=None, descriptions={ "errors": "workorder-shoc-emitter invocation errors", "throttles": "workorder-shoc-emitter invocation throttles", "duration": ( "workorder-shoc-emitter p99 duration approaching the 60s timeout" ), }, ) # --- Bespoke emitter alarms (plan Phase 4) --- # These don't fit add_standard_lambda_alarms' shape and stay bespoke. # Iterator age >= 10 min sustained means SHOC is likely down and the # shard is blocking (exactly the ordered-backpressure design working); # the two SQS-visible alarms page the operator replay runbook. shoc_emitter.metric( "IteratorAge", statistic="Maximum", period=Duration.minutes(5), ).create_alarm( stack, "ShocEmitterIteratorAgeAlarm", alarm_name="workorder-shoc-emitter-iterator-age", alarm_description=( "workorder-shoc-emitter stream lag >= 10 min " "(SHOC receiver likely down; shard blocking on retries)" ), threshold=600000, 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)) shoc_emitter_failures_queue.metric_approximate_number_of_messages_visible( period=Duration.minutes(5), statistic="Maximum", ).create_alarm( stack, "ShocEmitterFailuresMessagesAlarm", alarm_name="workorder-shoc-emitter-failures-messages", alarm_description=( "workorder-shoc-emitter retry-exhausted stream records parked " "(ESM failure metadata; replay rebuilds from DynamoDB)" ), 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)) shoc_emitter_rejected_queue.metric_approximate_number_of_messages_visible( period=Duration.minutes(5), statistic="Maximum", ).create_alarm( stack, "ShocEmitterRejectedMessagesAlarm", alarm_name="workorder-shoc-emitter-rejected-messages", alarm_description=( "workorder-shoc-emitter parked non-retryable 4xx deliveries " "(contract bug; inspect payloads and replay)" ), 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)) cdk.CfnOutput( stack, "ShocWebhookHmacSecretArn", value=shoc_hmac_secret.secret_arn, description=( "SHOC webhook HMAC secret ARN -- hand off to Luby (SHOC team) " "for the cross-account receiver fetch" ), )