From 13efee96825da02fd2759dbb18e5149214f69215 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri, 7 Aug 2026 13:42:35 -0400 Subject: [PATCH] feat(infra): add kebab WO tables and PLAT-11 cutover tooling (PLAT-11) (#172) * feat(infra): add kebab WO tables and PLAT-11 cutover tooling Create empty work-orders/work-order-comments under Terraform while writers stay on PascalCase; make emitter table classification env-driven and add same-account migrate/verify plus cutover runbook. * style(test): ruff-format shoc emitter envelope tests --- README.md | 14 +++ docs/plat-11/cutover-runbook.md | 92 ++++++++++++++ infra/shoc-assessment-reader/README.md | 10 +- lambdas/wo/shoc_emitter/envelope.py | 35 +++--- scripts/migrate_wo_tables.py | 166 +++++++++++++++++++++++++ scripts/replay_shoc_webhooks.py | 6 +- terraform/wo_ddb.tf | 47 +++++++ terraform/wo_shoc.tf | 3 + tests/test_shoc_emitter_envelope.py | 41 +++++- 9 files changed, 395 insertions(+), 19 deletions(-) create mode 100644 docs/plat-11/cutover-runbook.md create mode 100644 scripts/migrate_wo_tables.py diff --git a/README.md b/README.md index 8225ea7..ed71c20 100644 --- a/README.md +++ b/README.md @@ -272,6 +272,11 @@ The `purchase-orders` DynamoDB table is **owned by this repo's Terraform config* ### `WorkOrders` and `WorkOrderComments` tables (owned here) +> **PLAT-11 in progress.** Kebab physical names `work-orders` / `work-order-comments` +> are created empty under Terraform; writers/ESMs stay on PascalCase until the +> freeze cutover in [`docs/plat-11/cutover-runbook.md`](docs/plat-11/cutover-runbook.md). +> After cutover, this section will name the kebab tables. + Both tables are **owned by this repo's Terraform config** (`terraform/wo_ddb.tf`, retain lifecycle): - `WorkOrders` — PK `work_order_id` (S). @@ -394,6 +399,15 @@ python scripts/backfill_sites.py python scripts/replay_shoc_webhooks.py --url https://... --work-order-id 11144580730 # targeted, dry-run first python scripts/replay_shoc_webhooks.py --url https://... --work-order-id 11144580730 --execute # then re-POST python scripts/replay_shoc_webhooks.py --url https://... --since 2026-07-24T02:00:00Z --execute # everything touched since (full Scan) + +**PLAT-11 WO table rename** (PascalCase → kebab, freeze cutover). See +[`docs/plat-11/cutover-runbook.md`](docs/plat-11/cutover-runbook.md). Dry-run copy: + +```bash +python scripts/migrate_wo_tables.py copy +python scripts/migrate_wo_tables.py copy --execute +python scripts/migrate_wo_tables.py verify +``` ``` ## Directory Structure diff --git a/docs/plat-11/cutover-runbook.md b/docs/plat-11/cutover-runbook.md new file mode 100644 index 0000000..8547829 --- /dev/null +++ b/docs/plat-11/cutover-runbook.md @@ -0,0 +1,92 @@ +# PLAT-11 cutover runbook — WO Dynamo kebab rename + +Rename `WorkOrders` / `WorkOrderComments` → `work-orders` / `work-order-comments` +under HCP Terraform. SHOC is still **dev**; webhook quiet ~15–30 min is fine; +backfill after unfreeze. + +## Gates + +| Gate | Status / action | +|---|---| +| PLAT-86 soak ≥24h from Lambda touch `2026-08-07T15:06:49Z` | Ready at **2026-08-08T15:06:49Z** | +| Emitter failure/rejected SQS empty | Re-check before freeze | +| `shoc-assessment-dynamo-reader` | Torn down 2026-08-07 (prep) | +| SHOC Sync disable | Optional hygiene; not required | + +Freeze window: schedule after soak. Buffer ≤60 min. + +## Phase 1 — Prep apply (empty kebab tables) + +Already in tree: `aws_dynamodb_table.work_orders_kebab` / `work_order_comments_kebab`, +emitter env `WORK_ORDERS_TABLE` / `COMMENTS_TABLE` (still pointing at legacy), +`scripts/migrate_wo_tables.py`. + +1. Merge prep PR to `main`. +2. HCP Manual apply on `procurement-ingest-prod` (creates empty kebab tables only). +3. Confirm: + ```bash + aws dynamodb describe-table --table-name work-orders --profile seahaven-prod --region us-east-1 + aws dynamodb describe-table --table-name work-order-comments --profile seahaven-prod --region us-east-1 + ``` + +Writers and ESMs stay on PascalCase. + +## Phase 2 — Freeze cutover + +Record `FREEZE_START` (UTC ISO-8601) before step 1. + +1. **Freeze ingest** (pick one): + ```bash + aws lambda put-function-concurrency \ + --function-name workorder-email-processor \ + --reserved-concurrent-executions 0 \ + --profile seahaven-prod --region us-east-1 + ``` +2. **Disable ESMs** (note UUIDs from console or CLI list): + ```bash + aws lambda list-event-source-mappings \ + --function-name workorder-shoc-emitter \ + --profile seahaven-prod --region us-east-1 + # then for each UUID: + aws lambda update-event-source-mapping --uuid --enabled false \ + --profile seahaven-prod --region us-east-1 + ``` +3. **Copy**: + ```bash + python scripts/migrate_wo_tables.py copy --profile seahaven-prod + python scripts/migrate_wo_tables.py copy --profile seahaven-prod --execute + python scripts/migrate_wo_tables.py verify --profile seahaven-prod + ``` +4. **TF cutover** (same PR follow-up commit or cutover PR) — retarget: + - `wo_lambda.tf` / `api.tf` env → `aws_dynamodb_table.work_orders_kebab.name` (and comments kebab) + - `wo_shoc.tf` ESM `event_source_arn` → kebab stream ARNs; emitter env table names → kebab + - `iam.tf` DDB ARNs that reference `work_orders` / `work_order_comments` → include kebab resources (or switch) + - `wo_ddb.tf` `wo_ddb_alarm_tables` keys → `work-orders` / `work-order-comments` mapped to kebab resources +5. HCP Manual apply cutover. +6. **Unfreeze**: + ```bash + aws lambda delete-function-concurrency \ + --function-name workorder-email-processor \ + --profile seahaven-prod --region us-east-1 + # re-enable ESMs (or let TF `enabled = true` recreate/update) + ``` +7. Smoke: `scripts/post-deploy-smoke.sh`; SigV4 `GET /work-orders?limit=1`. +8. **Backfill SHOC**: + ```bash + python scripts/replay_shoc_webhooks.py \ + --url https://api.dev.seahaven.com/api/webhooks/work-orders \ + --since "$FREEZE_START" --execute + ``` + And/or one SHOC reconciliation pass against procurement-api. + +## Phase 3 — Decommission (after ≥24h on kebab) + +1. Lift `prevent_destroy` on **legacy** `work_orders` / `work_order_comments` only. +2. Remove legacy table resources + their import blocks + PascalCase alarm imports. +3. HCP apply (destroys PascalCase tables). +4. Update README + Confluence Architecture Map; note GitHub #24 superseded on PLAT-11. +5. Close PLAT-11 acceptance criteria. + +## Rollback (before decommission) + +Point env + ESMs back at PascalCase tables via TF; HCP apply. Legacy tables still hold data until Phase 3. diff --git a/infra/shoc-assessment-reader/README.md b/infra/shoc-assessment-reader/README.md index 659f053..d8dc0a5 100644 --- a/infra/shoc-assessment-reader/README.md +++ b/infra/shoc-assessment-reader/README.md @@ -1,4 +1,12 @@ -# shoc-assessment-dynamo-reader (TEMPORARY) +# shoc-assessment-dynamo-reader (TEMPORARY) — TORN DOWN + +**Status (2026-08-07):** Role deleted in seahaven-prod as part of PLAT-11 prep +(`./teardown-assessment-reader.sh`). Do not recreate for the kebab rename. +SHOC steady-state access remains `procurement-api` + webhook only. + +--- + +Original design (historical): Time-boxed cross-account read role in **seahaven-prod (011934824531)** that lets the Luby team assess the procurement-ingest DynamoDB tables directly from diff --git a/lambdas/wo/shoc_emitter/envelope.py b/lambdas/wo/shoc_emitter/envelope.py index 1552c48..c752244 100644 --- a/lambdas/wo/shoc_emitter/envelope.py +++ b/lambdas/wo/shoc_emitter/envelope.py @@ -1,17 +1,20 @@ """Stream-record -> SHOC webhook envelope mapping (contract sections 3-4). -Pure and total: no boto3 clients, no environment reads, no network. Every -function either returns a value or returns None (skip) -- a malformed stream -record must never raise here; classification gaps fall through to None so the -shard is never blocked by an unmappable record. ``build_event`` is the single -entry point the handler calls per record. +Pure and total for classification: no boto3 clients, no network. Table names +come from ``WORK_ORDERS_TABLE`` / ``COMMENTS_TABLE`` env (PLAT-11 kebab rename); +defaults preserve the pre-rename PascalCase names so unit tests without env +still pin the legacy stream ARNs. Every function either returns a value or +returns None (skip) -- a malformed stream record must never raise here; +classification gaps fall through to None so the shard is never blocked by an +unmappable record. ``build_event`` is the single entry point the handler calls +per record. Classification (docs/shoc-webhook-contract.md section 4): - WorkOrders INSERT -> work_order.created - WorkOrders MODIFY -> work_order.cancelled iff OLD wo_status was not - "cancelled" AND NEW wo_status is "cancelled", - else work_order.updated - WorkOrderComments INSERT -> work_order.comment_added + work-orders / WorkOrders INSERT -> work_order.created + work-orders / WorkOrders MODIFY -> work_order.cancelled iff OLD + wo_status was not "cancelled" AND NEW wo_status is "cancelled", + else work_order.updated + work-order-comments / WorkOrderComments INSERT -> work_order.comment_added REMOVE (either table) -> skip (no deletes in this pipeline) NEW write_origin == "shoc-write-api" -> skip (echo guard; a no-op branch today -- nothing writes that attribute yet -- kept so the phase-2 @@ -21,6 +24,7 @@ Classification (docs/shoc-webhook-contract.md section 4): Anything else (unknown table, comment MODIFYs) -> skip. """ +import os from datetime import datetime, timezone from decimal import Decimal @@ -30,8 +34,9 @@ SCHEMA_VERSION = 1 SOURCE = "procurement-ingest/workorder-shoc-emitter" CANCELLED_STATUS = "cancelled" SHOC_WRITE_ORIGIN = "shoc-write-api" -WORK_ORDERS_TABLE = "WorkOrders" -COMMENTS_TABLE = "WorkOrderComments" +# Defaults match pre-PLAT-11 physical names; TF/Lambda env overrides at cutover. +WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders") +COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments") # data payload field lists (contract sections 4.1 / 4.2). Absent attributes are # null-filled; source_email_s3_key is deliberately EXCLUDED (internal key, not @@ -100,9 +105,9 @@ def _table_name(event_source_arn: str) -> str | None: """Extract the exact table-name segment from a stream eventSourceARN. ARN format: arn:aws:dynamodb:region:acct:table/NAME/stream/LABEL. Parsing - the segment exactly (not a loose substring) matters: "WorkOrders" is a - prefix of nothing, but a substring test for it WOULD match - "WorkOrderComments"-adjacent names -- segment equality can't. + the segment exactly (not a loose substring) matters: legacy "WorkOrders" + is a leading substring of "WorkOrderComments", and kebab "work-orders" is + similarly adjacent to "work-order-comments" -- segment equality can't. """ parts = event_source_arn.split("/") if len(parts) >= _MIN_ARN_SLASH_SEGMENTS and parts[0].endswith(":table"): diff --git a/scripts/migrate_wo_tables.py b/scripts/migrate_wo_tables.py new file mode 100644 index 0000000..0139b9b --- /dev/null +++ b/scripts/migrate_wo_tables.py @@ -0,0 +1,166 @@ +#!/usr/bin/env python3 +"""PLAT-11 same-account WO table copy: PascalCase -> kebab-case. + +Copies ``WorkOrders`` -> ``work-orders`` and ``WorkOrderComments`` -> +``work-order-comments`` inside seahaven-prod. Dry-run is the DEFAULT; nothing +is written unless ``--execute`` is passed. + +Usage: + python scripts/migrate_wo_tables.py copy # dry-run both + python scripts/migrate_wo_tables.py copy --execute + python scripts/migrate_wo_tables.py verify + python scripts/migrate_wo_tables.py copy --table work-orders --execute +""" + +from __future__ import annotations + +import argparse +import sys +import time + +import boto3 +from boto3.dynamodb.types import TypeDeserializer + +EXPECTED_ACCOUNT = "011934824531" +REGION = "us-east-1" + +PAIRS = { + "work-orders": { + "source": "WorkOrders", + "dest": "work-orders", + "keys": ["work_order_id"], + }, + "work-order-comments": { + "source": "WorkOrderComments", + "dest": "work-order-comments", + "keys": ["work_order_id", "comment_id"], + }, +} + +_deserializer = TypeDeserializer() + + +def _client(profile: str): + session = boto3.Session(profile_name=profile, region_name=REGION) + acct = session.client("sts").get_caller_identity()["Account"] + if acct != EXPECTED_ACCOUNT: + sys.exit( + f"ERROR: profile {profile!r} resolves to account {acct}, " + f"expected {EXPECTED_ACCOUNT}. Aborting." + ) + return session.client("dynamodb") + + +def _describe_count(client, table: str) -> int: + return int(client.describe_table(TableName=table)["Table"]["ItemCount"]) + + +def _scan_all(client, table: str): + kwargs = {"TableName": table} + while True: + resp = client.scan(**kwargs) + for item in resp.get("Items", []): + yield item + if "LastEvaluatedKey" not in resp: + break + kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] + + +def _batch_write(client, table: str, items: list[dict], execute: bool) -> int: + if not execute: + return len(items) + written = 0 + for i in range(0, len(items), 25): + chunk = items[i : i + 25] + request_items = { + table: [{"PutRequest": {"Item": item}} for item in chunk], + } + unprocessed = request_items + while unprocessed: + resp = client.batch_write_item(RequestItems=unprocessed) + unprocessed = resp.get("UnprocessedItems") or {} + if unprocessed: + time.sleep(0.2) + written += len(chunk) + return written + + +def copy_table(client, pair_key: str, execute: bool) -> None: + pair = PAIRS[pair_key] + source, dest = pair["source"], pair["dest"] + print(f"==> copy {source} -> {dest} ({'EXECUTE' if execute else 'dry-run'})") + items = list(_scan_all(client, source)) + print(f" scanned {len(items)} items from {source}") + written = _batch_write(client, dest, items, execute=execute) + print(f" {'wrote' if execute else 'would write'} {written} items to {dest}") + + +def verify_table(client, pair_key: str) -> bool: + pair = PAIRS[pair_key] + source, dest = pair["source"], pair["dest"] + # ItemCount is eventually consistent; prefer live scan counts for gate. + src_items = list(_scan_all(client, source)) + dst_items = list(_scan_all(client, dest)) + src_n, dst_n = len(src_items), len(dst_items) + print(f"==> verify {source} ({src_n}) vs {dest} ({dst_n})") + if src_n != dst_n: + print(" FAIL count mismatch") + return False + + key_names = pair["keys"] + + def key_tuple(item): + return tuple(_deserializer.deserialize(item[k]) for k in key_names) + + src_by_key = {key_tuple(i): i for i in src_items} + dst_by_key = {key_tuple(i): i for i in dst_items} + missing = [k for k in src_by_key if k not in dst_by_key] + if missing: + print(f" FAIL missing keys on dest (showing up to 5): {missing[:5]}") + return False + + # Spot-check up to 25 items for full attribute equality. + checked = 0 + for key, src_item in list(src_by_key.items())[:25]: + if src_item != dst_by_key[key]: + print(f" FAIL attribute mismatch for key {key}") + return False + checked += 1 + print(f" OK counts equal; spot-checked {checked} items") + return True + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("command", choices=["copy", "verify"]) + parser.add_argument( + "--table", + choices=list(PAIRS.keys()), + help="Single pair key (default: both)", + ) + parser.add_argument( + "--execute", + action="store_true", + help="Actually write (copy only). Default is dry-run.", + ) + parser.add_argument("--profile", default="seahaven-prod") + args = parser.parse_args() + + client = _client(args.profile) + keys = [args.table] if args.table else list(PAIRS.keys()) + + if args.command == "copy": + for key in keys: + copy_table(client, key, execute=args.execute) + if not args.execute: + print("Dry-run complete. Re-run with --execute to write.") + return 0 + + ok = True + for key in keys: + ok = verify_table(client, key) and ok + return 0 if ok else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/replay_shoc_webhooks.py b/scripts/replay_shoc_webhooks.py index 66131ac..2e9bfdc 100644 --- a/scripts/replay_shoc_webhooks.py +++ b/scripts/replay_shoc_webhooks.py @@ -58,14 +58,16 @@ import urllib.request from datetime import datetime, timedelta, timezone from decimal import Decimal +import os + import boto3 from boto3.dynamodb.conditions import Key EXPECTED_ACCOUNT = "011934824531" REGION = "us-east-1" -WO_TABLE = "WorkOrders" -COMMENTS_TABLE = "WorkOrderComments" +WO_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders") +COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments") SECRET_NAME = "workorder-ingest/shoc-webhook-hmac" SCHEMA_VERSION = 1 diff --git a/terraform/wo_ddb.tf b/terraform/wo_ddb.tf index 536f530..337560c 100644 --- a/terraform/wo_ddb.tf +++ b/terraform/wo_ddb.tf @@ -1,3 +1,7 @@ +# PLAT-11: PascalCase tables remain authoritative until the freeze cutover. +# Kebab tables are created empty in the prep apply; writers/ESMs stay on legacy +# until docs/plat-11/cutover-runbook.md Phase 2. + resource "aws_dynamodb_table" "work_orders" { name = "WorkOrders" billing_mode = "PAY_PER_REQUEST" @@ -40,7 +44,50 @@ resource "aws_dynamodb_table" "work_order_comments" { } } +resource "aws_dynamodb_table" "work_orders_kebab" { + name = "work-orders" + billing_mode = "PAY_PER_REQUEST" + hash_key = "work_order_id" + + attribute { + name = "work_order_id" + type = "S" + } + + stream_enabled = true + stream_view_type = "NEW_AND_OLD_IMAGES" + + lifecycle { + prevent_destroy = true + } +} + +resource "aws_dynamodb_table" "work_order_comments_kebab" { + name = "work-order-comments" + billing_mode = "PAY_PER_REQUEST" + hash_key = "work_order_id" + range_key = "comment_id" + + attribute { + name = "work_order_id" + type = "S" + } + + attribute { + name = "comment_id" + type = "S" + } + + stream_enabled = true + stream_view_type = "NEW_AND_OLD_IMAGES" + + lifecycle { + prevent_destroy = true + } +} + locals { + # Alarms stay on the live writer tables (legacy until cutover). wo_ddb_alarm_tables = { "WorkOrders" = aws_dynamodb_table.work_orders.name "WorkOrderComments" = aws_dynamodb_table.work_order_comments.name diff --git a/terraform/wo_shoc.tf b/terraform/wo_shoc.tf index d576b6f..dee5526 100644 --- a/terraform/wo_shoc.tf +++ b/terraform/wo_shoc.tf @@ -292,6 +292,9 @@ resource "aws_lambda_function" "wo_shoc_emitter" { HMAC_SECRET_ARN = var.shoc_hmac_secret_arn REJECTED_QUEUE_URL = aws_sqs_queue.shoc_emitter_rejected.url SHOC_WEBHOOK_URL = var.shoc_webhook_url + # Stream ARN classification (PLAT-11); must match ESM source table names. + WORK_ORDERS_TABLE = aws_dynamodb_table.work_orders.name + COMMENTS_TABLE = aws_dynamodb_table.work_order_comments.name } } diff --git a/tests/test_shoc_emitter_envelope.py b/tests/test_shoc_emitter_envelope.py index 7ce2e22..50ebdf1 100644 --- a/tests/test_shoc_emitter_envelope.py +++ b/tests/test_shoc_emitter_envelope.py @@ -12,7 +12,8 @@ contract names AND the published spec's webhooks keys (lambdas/api/ openapi.json), and the classification's status/record_type value space must stay anchored to the WO pipeline's template_parser enums. -Offline and pure -- envelope.py has no boto3 clients, no env reads. +Offline -- envelope.py has no boto3 clients and no network. Table names +default to PascalCase; set WORK_ORDERS_TABLE / COMMENTS_TABLE for kebab. """ import json @@ -231,6 +232,44 @@ def test_table_discrimination_is_exact_segment_not_substring(envelope): assert envelope.build_event(_record("", "INSERT", _wo_image())) is None +def test_kebab_table_names_via_env(monkeypatch): + """PLAT-11: classification follows WORK_ORDERS_TABLE / COMMENTS_TABLE.""" + import sys + + monkeypatch.setenv("WORK_ORDERS_TABLE", "work-orders") + monkeypatch.setenv("COMMENTS_TABLE", "work-order-comments") + sys.modules.pop("wo_shoc_emitter_envelope", None) + kebab_envelope = load_lambda_module("wo", "shoc_emitter/envelope") + wo_arn = ( + "arn:aws:dynamodb:us-east-1:011934824531:table/work-orders" + "/stream/2026-08-07T00:00:00.000" + ) + comments_arn = ( + "arn:aws:dynamodb:us-east-1:011934824531:table/work-order-comments" + "/stream/2026-08-07T00:00:00.000" + ) + assert ( + kebab_envelope.build_event(_record(wo_arn, "INSERT", _wo_image()))["event_type"] + == "work_order.created" + ) + assert ( + kebab_envelope.build_event(_record(comments_arn, "INSERT", _comment_image()))[ + "event_type" + ] + == "work_order.comment_added" + ) + # Legacy PascalCase ARNs must not match kebab env names. + assert ( + kebab_envelope.build_event(_record(WO_STREAM_ARN, "INSERT", _wo_image())) + is None + ) + # Restore cached module for later tests that expect PascalCase defaults. + sys.modules.pop("wo_shoc_emitter_envelope", None) + monkeypatch.delenv("WORK_ORDERS_TABLE", raising=False) + monkeypatch.delenv("COMMENTS_TABLE", raising=False) + load_lambda_module("wo", "shoc_emitter/envelope") + + # --- Envelope fields ---------------------------------------------------------