diff --git a/README.md b/README.md index 2628318..7ee64b2 100644 --- a/README.md +++ b/README.md @@ -91,7 +91,7 @@ Amazon APM work order emails (from Hexagon EAM / HxGN SmartCloud) are received a **Flow:** `WorkOrders` / `WorkOrderComments` DynamoDB Streams (`NEW_AND_OLD_IMAGES` — the OLD image is what lets the emitter detect the `wo_status → cancelled` transition) → `workorder-shoc-emitter` → HMAC-signed HTTPS POST → SHOC backend. Every WO mutation becomes one webhook event (`work_order.created` / `.updated` / `.cancelled` / `.comment_added`) within seconds of the DynamoDB commit; DynamoDB stays the source of truth. - **Contract:** `docs/shoc-webhook-contract.md` (Rev 2026-07-23) is the producer/consumer contract SHOC builds its receiver against, and `docs/shoc-webhook-test-vectors.json` is the shared receiver-verification vector set — both sides pin their HMAC implementation against the same vectors (producer-side via the golden-vector tests over `delivery.sign_body`). -- **Ships DARK.** Both event-source mappings deploy with `enabled=False`: the full stack (Lambda, queues, alarms, secret, rotation) deploys and is testable with zero deliveries while SHOC has no receiver, so this repo's merge cadence never depends on SHOC's. Activation is a deliberate one-line `enabled=True` PR, gated on the SHOC receiver passing the shared test vectors. The ESMs start at `LATEST` — no historical flood; SHOC loads history through the `procurement-api` read API at activation instead. +- **ACTIVE since 2026-07-30.** Both event-source mappings run with `enabled=True`, flipped after the SHOC receiver on `api.dev.seahaven.com` passed the shared test vectors live (valid current-kid signature accepted, duplicate `delivery_id` deduplicated, tampered/stale/unknown-kid all rejected 401). The stack originally shipped dark (`enabled=False`) so it could deploy and be tested with zero deliveries while SHOC had no receiver. The ESMs start at `LATEST` — no historical flood; SHOC backfills history through the `procurement-api` read API, not the stream. - **Ordering/retry semantics.** `parallelization_factor=1`, `bisect_batch_on_error=False`, `retry_attempts=-1`, `maximum_record_age=24h`: a retryable failure (429/5xx/timeout/connection error) blocks the shard and retries from the failed record — per-work-order commit order is the guarantee, and blocking is the intended behavior when SHOC is down. `report_batch_item_failures` keeps earlier in-batch successes from being re-delivered. Records that exhaust the 24h age are parked as ESM **failure metadata** (not full records) on `workorder-shoc-emitter-failures` (`on_failure` destination); a non-retryable 4xx (a contract bug, never worth blocking the shard for 24h) parks the **full `{envelope, response_status}` payload** on `workorder-shoc-emitter-rejected` and the loop continues. Both queues: 14-day retention, SSL-enforced, alarmed (see [CloudWatch alarms](#cloudwatch-alarms)); recovery is `scripts/replay_shoc_webhooks.py` (see [Scripts](#scripts)). - **Secret + KMS.** The HMAC signing keys live in Secrets Manager secret `workorder-ingest/shoc-webhook-hmac` (value `{"keys": [{"kid", "secret"}, ...]}`, newest first, max 2), encrypted with the dedicated CMK `workorder-ingest-shoc-webhook-kms` — deliberately **not** `alias/seahaven-dynamodb`, so the SHOC cross-account grant's decrypt reach covers exactly this one secret and nothing else. The secret's removal policy is **`DESTROY`, deliberately not `RETAIN`**: the value is machine-generated HMAC material, fully regenerable by a single rotation, so `RETAIN` buys nothing and would expose the fixed-name RETAIN-orphan deadlock (a failed create orphans an empty shell holding the global name; every later create fails `AlreadyExists`). **Accidental-deletion recovery runbook:** redeploy to recreate the secret, force a rotation (`aws secretsmanager rotate-secret --secret-id workorder-ingest/shoc-webhook-hmac`), notify the SHOC team — receivers re-fetch within their ≤5-minute cache TTL, so no coordination window is needed — then watch the `-failures` queue and replay the gap with the replay script. - **Rotation.** `workorder-shoc-hmac-rotator` runs on a 30-day schedule: it prepends a fresh 64-hex-char key as `keys[0]` and truncates the list to 2 entries (one overlap cycle). `kid` format is `YYYY-MM-DDTHH`. The emitter always signs with `keys[0]` behind a 5-minute TTL cache; the receiver accepts any listed `kid` and re-fetches on an unknown one — there is no delivery window in which signatures can't verify. @@ -317,7 +317,7 @@ Both tables are **owned by this repo's `WorkorderIngestStack`** (`cdk/wo_stack.p Both currently use default DynamoDB encryption — they are **not** yet on the shared customer-managed CMK (`alias/seahaven-dynamodb`, INFRA-95 / M-3); that migration is tracked in INFRA-6. The `workorder-email-processor` role no longer holds a pre-emptive encrypt/decrypt grant on that CMK (removed in the 2026-06-17 security sweep — it was unused while the tables are unencrypted and extended the role's decrypt reach to the CMK protecting `purchase-orders`). Re-add the grant as part of the INFRA-6 migration, at which point `grant_read_write_data` on the then-encrypted tables propagates the needed key permissions automatically. -**Consumers (read-only) — data contract:** `seahaven-slack-bot` (the former external reader) was decommissioned 2026-07-23; its grants are gone. Current consumers: the `procurement-api` Lambda (this repo, `Table.from_table_name` + `grant_read_data`, serving `GET /work-orders*`), and the `workorder-shoc-emitter` stream consumer (this repo) — both tables now stream `NEW_AND_OLD_IMAGES`, and the emitter is a live consumer of those streams, dark (ESMs disabled) until the SHOC activation PR. External readers (SHOC) consume through the REST contract (`lambdas/api/openapi.json`) and the webhook contract (`docs/shoc-webhook-contract.md`); both documents' field lists mirror `lambdas/wo/email_processor/persistence.py`. Any change to table name, key schema, attribute names, or encryption configuration (e.g. the INFRA-6 CMK migration) must update those two contracts in the same PR — the tables are imported by name, so there is no compile-time link and breakage surfaces at runtime. +**Consumers (read-only) — data contract:** `seahaven-slack-bot` (the former external reader) was decommissioned 2026-07-23; its grants are gone. Current consumers: the `procurement-api` Lambda (this repo, `Table.from_table_name` + `grant_read_data`, serving `GET /work-orders*`), and the `workorder-shoc-emitter` stream consumer (this repo) — both tables now stream `NEW_AND_OLD_IMAGES`, and the emitter is an active consumer of those streams (ESMs enabled 2026-07-30). External readers (SHOC) consume through the REST contract (`lambdas/api/openapi.json`) and the webhook contract (`docs/shoc-webhook-contract.md`); both documents' field lists mirror `lambdas/wo/email_processor/persistence.py`. Any change to table name, key schema, attribute names, or encryption configuration (e.g. the INFRA-6 CMK migration) must update those two contracts in the same PR — the tables are imported by name, so there is no compile-time link and breakage surfaces at runtime. **Stream field contract (`site_code`).** Three definitions of "is this a valid site code" have existed in this repo at once. The canonical shape is `derived_fields._STRICT_CODE_RE` (`[A-Z][A-Z0-9]{2,4}`, `fullmatch`) plus its skip-list semantics (a code-shaped token is only a real site code if it is *not* skip-listed, e.g. `LLC`/`INC`/`CORP`/`LTD`/`ATTN`), used by `derive_site_code()` (the `enrich_parsed()` classifier, see the **Derived fields** note in the Purchase Orders flow above). Honestly noted, not papered over: `lambdas/po/site_extractor/handler.py`'s own `SITE_CODE_PATTERN` (`[A-Z]{2,4}\d{1,2}`, prefix-anchored `.match`, digit-requiring) still diverges from the canonical shape as of this phase — it rejects valid all-letter codes like `KLAL` (which surface as permanent `pending-site-review` rows) and accepts overlong junk like `DLI6X`/`SNY55` that the canonical `fullmatch` would not. Reconciling `po-ingest-site-extractor`'s direct-field validation onto the canonical `derived_fields` shape is Phase 6 scope, tracked separately — this paragraph will be updated when it lands. @@ -508,7 +508,7 @@ lambdas/ # Phase 2: shared Code.from_asset("../lambdas") bundling ses-stamped/ auth-pass-01.eml # Phase 8: the one new fixture allowed this phase -- a synthesized Authentication-Results header block (from WO_SES_HEADER in test_ses_auth.py), not scraped mail web_ui/ # handler.py imports `from web_ui_auth import is_authenticated` (shared) - shoc_emitter/ # SHOC webhook emitter (ships dark -- ESMs disabled; see the SHOC + shoc_emitter/ # SHOC webhook emitter (ACTIVE since 2026-07-30; see the SHOC # webhook feed section). Flat siblings, bare-name imports handler.py # thin per-record event loop + partial-batch failure report + rejected-queue parking envelope.py # pure stream-record -> envelope mapping + event classification (incl. echo guard) diff --git a/cdk/wo_stack.py b/cdk/wo_stack.py index 7d6bbc0..bb133c5 100644 --- a/cdk/wo_stack.py +++ b/cdk/wo_stack.py @@ -1,29 +1,56 @@ """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, - aws_ses as ses, - aws_ses_actions as ses_actions, +) +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 -import common - class WorkorderIngestStack(Stack): def __init__(self, scope: Construct, construct_id: str, **kwargs): @@ -750,14 +777,12 @@ def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_to shoc_emitter_rejected_queue.grant_send_messages(shoc_emitter) # ------------------------------------------------------------------ - # SHIPS DARK: both event source mappings deploy with enabled=False, - # deliberately. The full stack (Lambda, queues, alarms, secret, - # rotation) deploys and is testable with zero deliveries while SHOC - # has no receiver, so our merge cadence never depends on Luby's. - # Activation is a one-line enabled=True PR gated on the SHOC receiver - # passing the shared HMAC test vectors (plan Phase 3). LATEST start - # position => no historical flood at activation; SHOC loads history - # via the procurement read API instead. + # 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 @@ -775,7 +800,7 @@ def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_to max_record_age=Duration.hours(24), parallelization_factor=1, report_batch_item_failures=True, - enabled=False, + enabled=True, on_failure=lambda_event_sources.SqsDlq(shoc_emitter_failures_queue), ) ) @@ -789,7 +814,7 @@ def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_to max_record_age=Duration.hours(24), parallelization_factor=1, report_batch_item_failures=True, - enabled=False, + enabled=True, on_failure=lambda_event_sources.SqsDlq(shoc_emitter_failures_queue), ) )