From 073201f633d50cee539ff213f2078e9424bead33 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 23 Jul 2026 17:08:47 -0400 Subject: [PATCH] Migrate to seahaven-prod: deploy role, backfill tooling, account-portability fixes (#125) * feat(migration): prepare stacks and tooling for the seahaven-prod account move Phase 1 of the mgmt (328440206208) -> seahaven-prod (011934824531) migration. No behavior change in-account; everything here is additive or account-portability hygiene: - infra/deploy-role/: reviewed OIDC deploy-role artifacts for prod (trust main-only, cdk-hnb659fds-* AssumeRole, smoke-invoke-lambda scoped to exactly the two email-processor fn ARNs). Codifies the previously out-of-band smoke-invoke grant. - Table resource policies: make_slack_bot_read_policy in cdk/common.py, applied to purchase-orders, verified-sites, WorkOrders, WorkOrderComments (NOT pending-site-review; no bot consumer). Grants the mgmt-resident seahaven-slack-bot roles read-only cross-account access post-move (bot-side identity grants land in the slack-bot repo). - scripts/migrate_tables.py: dry-run-default backfill tool implementing the plan's per-table semantics (superset overwrite, ingested_at cutoff for WorkOrderComments, backup-gated truncate-and-load for the two site tables) plus a verify subcommand (count parity, spot checks, sticky-Cancelled drift check). - tests/test_resource_policy_helper.py: statement-shape unit tests + static pins that exactly the four bot-read tables carry the policy. - Account-literal fixes: account-agnostic fixture bucket in test_reprocess_contract; runbook/README/po-template-parser account references updated to prod with historical mgmt notes; README gains the account-prerequisites list (imported-by-name dependencies). deploy.yaml is deliberately unchanged (push-to-main auto-deploy kept). Merge is held until migration Phase 0 completes; flipping the AWS_DEPLOY_ROLE_ARN repo secret and merging this PR IS the first prod deploy. * fix(migration): verify backup AVAILABLE pre-truncate; document wildcard risk acceptance (cross-review FIX/NIT) * refactor(migration): drop cross-account read grants (slack-bot decommissioned); harden backfill + deploy role seahaven-slack-bot was decommissioned 2026-07-23 (stack DELETE_IN_PROGRESS, consumer Lambdas gone); its successor sh-mcp is undeployed and uses same-account DynamoDB access. So no live consumer reads these tables cross-account. Per Adam's call, drop the cross-account grants entirely and re-add correctly-scoped ones if/when sh-mcp deploys to a different account. - Remove the four table resource policies + make_slack_bot_read_policy helper + its constants (cdk/common.py, po_stack.py, wo_stack.py) and the helper's unit test. Both stacks synth with zero table ResourcePolicy. - scripts/migrate_tables.py hardening (fixes from the sh-security-review fan-out on the destructive backfill tool): * validate --cutoff strictly (parse ISO-8601, require aware UTC, re-emit canonical second-precision form) so a malformed cutoff can't silently copy dual-window rows or drop history; * reject `copy --all` up front (must run tables individually, in order, with the stream-drain wait) instead of writing three tables then erroring; * truncate backup gate now also checks recency (<1h) and TableId, not just status+name; * verify requires --cutoff whenever a cutoff table is in scope (else it false-flags dual-window rows as MISSING); * sticky-cancel is now PREVENTED copy-side (a non-Cancelled source item never overwrites a dest-Cancelled PO), and the verify comment no longer overstates what its source-side scan covers; * spot-check all modes (truncate_load keys are verbatim, so key-existence is sound there too). - Deploy role: scope cloudformation:DescribeStacks to this repo's stacks + CDKToolkit (was Resource:*, disclosed all tenant stacks in the shared prod account); add a drift check warning on unexpected role policies and drop the dead SMOKE_POLICY_NAME var; document the shared-account bootstrap-role accepted risk in the deploy-role README. * docs(deploy-role): fold in cross-review NITs (DescribeStacks maintenance note, warn-only drift rationale) * ci: update workflow to use new workflow tag (ruff versioning fix) * fix(migration): address Open SWE review findings on migrate_tables.py - Validate the truncate backup on dry-run as well as --execute so a missing/stale/wrong-incarnation --backup-arn surfaces on the rehearsal run (finding f_24a48b8900). - Assert configured keys match the live key schema of both tables before any key projection, turning config/schema drift into a descriptive abort instead of a mid-backfill KeyError (finding f_cb6b5a6c59). - Clarify why key-existence spot-checks are sound for WorkOrderComments: the copy Puts source items verbatim and the sample uses the same cutoff filter, so per-account comment_id divergence never enters the check (finding f_390b7d6c3b is a false positive; comment hardened). --- .github/workflows/ci.yaml | 2 +- README.md | 16 +- docs/po-template-parser.md | 6 +- docs/runbook-dlq-recovery.md | 12 +- infra/deploy-role/README.md | 34 ++ infra/deploy-role/create-deploy-role.sh | 96 +++++ infra/deploy-role/permissions-policy.json | 40 +++ infra/deploy-role/trust-policy.json | 18 + scripts/migrate_tables.py | 416 ++++++++++++++++++++++ tests/test_reprocess_contract.py | 4 +- 10 files changed, 636 insertions(+), 8 deletions(-) create mode 100644 infra/deploy-role/README.md create mode 100755 infra/deploy-role/create-deploy-role.sh create mode 100644 infra/deploy-role/permissions-policy.json create mode 100644 infra/deploy-role/trust-policy.json create mode 100644 scripts/migrate_tables.py diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 71acfcb..41319cb 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -8,7 +8,7 @@ permissions: jobs: ci: - uses: Sea-Haven-Industries/.github/.github/workflows/ci-python-sam.yaml@fd60e4c9041784f666ac0fdefb9bec3c7fbf5143 # main + uses: Sea-Haven-Industries/.github/.github/workflows/ci-python-sam.yaml@f71002a9ed2938730b683249b28059c92a081af6 with: source-dirs: "lambdas cdk tests scripts" run-tests: true diff --git a/README.md b/README.md index 0067b50..d55ab78 100644 --- a/README.md +++ b/README.md @@ -88,7 +88,7 @@ Amazon APM work order emails (from Hexagon EAM / HxGN SmartCloud) are received a ## Architecture -**IaC:** AWS CDK (Python), two stacks in one app, region `us-east-1`. The `cdk.Environment` is deliberately **account-agnostic** (region-only, no `account=`): the stacks deploy to whichever account the deploy credentials target (`328440206208` today), and every account-derived template value — bucket names, `Lambda::Permission` source account, the site-alerts SNS action ARN, the Bedrock ARN below — renders as the CloudFormation `AWS::AccountId` pseudo-parameter rather than a literal. Pinning `account=` was evaluated in Phase 4 and rejected: it would resolve those tokens to literals, and against the deployed (account-agnostic) templates CloudFormation flags the `RemovalPolicy.RETAIN` email buckets as requiring replacement — a data-loss risk — for no functional gain. +**IaC:** AWS CDK (Python), two stacks in one app, region `us-east-1`. The `cdk.Environment` is deliberately **account-agnostic** (region-only, no `account=`): the stacks deploy to whichever account the deploy credentials target (`011934824531` seahaven-prod since the 2026-07 account migration; formerly mgmt `328440206208`), and every account-derived template value — bucket names, `Lambda::Permission` source account, the site-alerts SNS action ARN, the Bedrock ARN below — renders as the CloudFormation `AWS::AccountId` pseudo-parameter rather than a literal. Pinning `account=` was evaluated in Phase 4 and rejected: it would resolve those tokens to literals, and against the deployed (account-agnostic) templates CloudFormation flags the `RemovalPolicy.RETAIN` email buckets as requiring replacement — a data-loss risk — for no functional gain. All Lambdas: Python 3.12, ARM64, 60-day log retention. @@ -313,6 +313,20 @@ Branch protection on `main` — all changes through PR. ## Setup +**Account prerequisites** — the stacks import four dependencies by name, so all +of these must exist in the target account BEFORE the first deploy (in +seahaven-prod they are provisioned by the seahaven-org-baseline repo and the +migration Phase 0 runbook): + +- SES receipt rule set `INBOUND_MAIL` (may be inactive; the stacks attach + their receipt rules to it) + a verified `int.seahaven.com` domain identity +- SNS topic `site-alerts` (with the `alias/seahaven-alarm-topics` CMK) +- SSM param `/seahaven/dynamodb/cmk-arn` -> KMS `alias/seahaven-dynamodb` +- Secrets Manager secret `procurement-ingest/web-ui-auth-token` (step 3) +- OIDC deploy role `githubdeploy-procurement-ingest` (see `infra/deploy-role/`) + with the `smoke-invoke-lambda` policy, or the post-deploy smoke gate fails + closed + 1. Bootstrap CDK: `cdk bootstrap aws://{AccountId}/us-east-1` 2. Ensure the Bedrock inference profile `us.anthropic.claude-haiku-4-5-20251001-v1:0` is enabled in `us-east-1` (it is; the CDK grants cover cross-region routing to `us-east-2`/`us-west-2`). No API key or secret to set — the processors authenticate to Bedrock via their IAM roles. 3. Create the web UI auth-gate shared secret. This secret is **imported by name** diff --git a/docs/po-template-parser.md b/docs/po-template-parser.md index 7c087e2..193677f 100644 --- a/docs/po-template-parser.md +++ b/docs/po-template-parser.md @@ -29,6 +29,10 @@ The PO email processor (`lambdas/po/email_processor/handler.py`) sends **every** Two phases: a read-only multi-agent (ultracode) feasibility study over a 120-email sample, then a full-bucket triage over **all 3,448** inbound emails to get real distribution numbers. ### 2.1 Data access +> **Migration note (2026-07):** this section is historical — the harvest ran +> against the management account. Post-migration the live bucket is +> `s3://po-ingest-emails-011934824531/inbound/` (seahaven-prod, profile +> `seahaven-prod`); the mgmt bucket persists only until decommission. - PO email bucket: `s3://po-ingest-emails-328440206208/inbound/` (AWS account **328440206208**, us-east-1). - ~~Reached via AWS profile `amoussa-mgmt` (SSO); default CLI creds are the personal account `681986854588`~~ **Stale (corrected 2026-07-16 during the fixture harvest):** the default CLI session is now authenticated to **328440206208** directly — no `--profile` flag needed. - ⚠️ **Corpus is aging out:** `inbound/` objects carry an S3 lifecycle expiration (~90-day rolling window; oldest object 2026-04-17 at harvest time, 3,422 objects vs 3,448 at triage). Any further harvesting should not be deferred long. @@ -279,4 +283,4 @@ Plus `missing_required_field` for the required labeled fields (per-item metadata - PO handler + `EXTRACTION_PROMPT`: `lambdas/po/email_processor/handler.py`. - PO stack (alarms, IAM, Bedrock): `cdk/po_stack.py`. - Bedrock model: inference profile `us.anthropic.claude-haiku-4-5-20251001-v1:0`. -- PO email bucket: `s3://po-ingest-emails-328440206208/inbound/` (account 328440206208, profile `amoussa-mgmt`). +- PO email bucket (historical, pre-migration): `s3://po-ingest-emails-328440206208/inbound/`; post-migration `s3://po-ingest-emails-011934824531/inbound/` (seahaven-prod). diff --git a/docs/runbook-dlq-recovery.md b/docs/runbook-dlq-recovery.md index d945f03..31331a6 100644 --- a/docs/runbook-dlq-recovery.md +++ b/docs/runbook-dlq-recovery.md @@ -18,14 +18,18 @@ There is **no console redrive-to-source**: the SQS console's "redrive to source" applies only to SQS-to-SQS DLQ relationships, not to a Lambda `DeadLetterConfig` target. Recovery is manual, via targeted re-invoke. -## Resource inventory (acct 328440206208, us-east-1) +## Resource inventory (acct 011934824531 seahaven-prod, us-east-1) | Pipeline | Function | DLQ queue name | Raw-email bucket | DLQ alarm | |---|---|---|---|---| -| PO | `po-email-processor` | `po-ingest-EmailProcessorDlqA753DED5-az8LUZE3ubtz` | `po-ingest-emails-328440206208` | `po-email-processor-dlq-messages` | -| WO | `workorder-email-processor` | `WorkorderIngestStack-EmailProcessorDlqA753DED5-Q8H555LrqSU1` | `workorder-ingest-emails-328440206208` | `workorder-email-processor-dlq-messages` | +| PO | `po-email-processor` | _CDK-generated; fill from stack resources after the first prod deploy_ | `po-ingest-emails-011934824531` | `po-email-processor-dlq-messages` | +| WO | `workorder-email-processor` | _CDK-generated; fill from stack resources after the first prod deploy_ | `workorder-ingest-emails-011934824531` | `workorder-email-processor-dlq-messages` | -DLQ URLs are `https://sqs.us-east-1.amazonaws.com/328440206208/`. +DLQ URLs are `https://sqs.us-east-1.amazonaws.com/011934824531/`. +(Until mgmt decommission completes, the pre-migration mgmt-account queues +`po-ingest-EmailProcessorDlqA753DED5-az8LUZE3ubtz` and +`WorkorderIngestStack-EmailProcessorDlqA753DED5-Q8H555LrqSU1` still exist in +328440206208 with any pre-cutover dead letters.) Both DLQs: 14-day retention, SSE, TLS-enforced, `VisibilityTimeout` 30s. ## Recovery procedure (no redrive — receive → extract key → targeted re-invoke → verify → purge) diff --git a/infra/deploy-role/README.md b/infra/deploy-role/README.md new file mode 100644 index 0000000..5dbfec5 --- /dev/null +++ b/infra/deploy-role/README.md @@ -0,0 +1,34 @@ +# Deploy role: githubdeploy-procurement-ingest (seahaven-prod) + +OIDC deploy role for this repo's GitHub Actions pipeline in AWS account +`011934824531` (seahaven-prod), us-east-1. Created as part of the migration +from the management account (328440206208); the mgmt role of the same name +stays untouched until decommission as the emergency mgmt deploy path. + +## Files + +| File | Purpose | +|---|---| +| `trust-policy.json` | OIDC trust: `repo:Sea-Haven-Industries/procurement-ingest:ref:refs/heads/main` only | +| `permissions-policy.json` | `sts:AssumeRole` on the four `cdk-hnb659fds-*` bootstrap roles, `cloudformation:DescribeStacks` scoped to this repo's stacks + `CDKToolkit` (cd-cdk health check), and `lambda:InvokeFunction` on exactly the two email-processor function ARNs (post-deploy smoke gate) | +| `create-deploy-role.sh` | Idempotent create-or-update from the two JSON files, profile `seahaven-prod` | + +> **Maintenance note:** `DescribeStacks` is scoped to `stack/po-ingest/*`, `stack/WorkorderIngestStack/*`, and `stack/CDKToolkit/*`. If a third stack is ever added to this CDK app, add its ARN pattern here and re-run the review-then-apply flow — otherwise the cd-cdk health check on the new stack will `AccessDenied`. + +## Why the SmokeInvokeLambda statement exists + +`deploy.yaml` runs `scripts/post-deploy-smoke.sh` under the deploy role's own +session, not the assumed `cdk-*` roles. Without `lambda:InvokeFunction` on the +two function ARNs the smoke gate hits AccessDenied and every deploy fails +closed. The mgmt-era grant was applied out-of-band and undocumented; keeping +it in these reviewed artifacts closes that gap. Scope it to exactly the two +ARNs, never `Resource: "*"`. + +## Change process + +1. Edit the JSON artifacts on a branch; both gates must pass on the exact + files before anything is applied: GPT-4.1 cross-family review + (cross_review.py) and /sh-security-review. +2. Run `./create-deploy-role.sh` (idempotent) with the seahaven-prod profile. +3. Verify: `aws iam simulate-principal-policy` for the bootstrap-role + AssumeRole and both InvokeFunction ARNs, then a real pipeline run. diff --git a/infra/deploy-role/create-deploy-role.sh b/infra/deploy-role/create-deploy-role.sh new file mode 100755 index 0000000..ee2587b --- /dev/null +++ b/infra/deploy-role/create-deploy-role.sh @@ -0,0 +1,96 @@ +#!/usr/bin/env bash +############################################################################### +# create-deploy-role.sh +# +# Creates / updates the GitHub Actions OIDC deploy role +# `githubdeploy-procurement-ingest` in AWS account 011934824531 (seahaven-prod), +# us-east-1, for the Sea-Haven-Industries/procurement-ingest repo (main branch). +# +# Part of the mgmt -> seahaven-prod migration. The mgmt-account role of the +# same name is left untouched until decommission (emergency mgmt deploy path). +# +# GATE - DO NOT EXECUTE until BOTH of the following have passed on these +# exact artifact files: +# 1. GPT-4.1 cross-family review (IAM policy / trust-policy change) via +# cross_review.py in the security-review repo. +# 2. /sh-security-review (deep agentic pass - IaC/IAM is a gated surface). +# +# The SmokeInvokeLambda statement is REQUIRED, not optional: deploy.yaml runs +# scripts/post-deploy-smoke.sh under this role's session (not the assumed +# cdk-* roles), and the smoke gate fails closed without lambda:InvokeFunction +# on exactly the two email-processor function ARNs. The mgmt-era equivalent +# of this policy was applied out-of-band and documented nowhere; these +# artifacts are the fix for that debt. +# +# Resolved facts (verified 2026-07-23, read-only): +# Account : 011934824531 (seahaven-prod) +# Region : us-east-1 +# Qualifier : hnb659fds (AWS CDK DEFAULT - matches cdk/cdk.json) +# OIDC prov : arn:aws:iam::011934824531:oidc-provider/token.actions.githubusercontent.com +############################################################################### + +set -euo pipefail + +PROFILE="seahaven-prod" +ROLE_NAME="githubdeploy-procurement-ingest" +POLICY_NAME="procurement-ingest-deploy" +ACCOUNT_ID="011934824531" +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +TRUST_POLICY="file://${SCRIPT_DIR}/trust-policy.json" +PERMS_POLICY="file://${SCRIPT_DIR}/permissions-policy.json" + +echo "==> Verifying active account for profile '${PROFILE}'..." +CALLER_ACCOUNT="$(aws --profile "${PROFILE}" sts get-caller-identity --query Account --output text)" +if [[ "${CALLER_ACCOUNT}" != "${ACCOUNT_ID}" ]]; then + echo "ERROR: profile '${PROFILE}' resolves to account ${CALLER_ACCOUNT}, expected ${ACCOUNT_ID}. Aborting." >&2 + exit 1 +fi + +echo "==> Ensuring role '${ROLE_NAME}' exists with the correct trust policy..." +if aws --profile "${PROFILE}" iam get-role --role-name "${ROLE_NAME}" >/dev/null 2>&1; then + echo " Role exists - updating assume-role (trust) policy." + aws --profile "${PROFILE}" iam update-assume-role-policy \ + --role-name "${ROLE_NAME}" \ + --policy-document "${TRUST_POLICY}" +else + echo " Role absent - creating." + aws --profile "${PROFILE}" iam create-role \ + --role-name "${ROLE_NAME}" \ + --assume-role-policy-document "${TRUST_POLICY}" \ + --description "GitHub Actions OIDC deploy role for Sea-Haven-Industries/procurement-ingest (main)" \ + --max-session-duration 3600 \ + --tags Key=project,Value=procurement-ingest Key=managed-by,Value=create-deploy-role.sh +fi + +echo "==> Putting inline permissions policy '${POLICY_NAME}' (create-or-replace)..." +aws --profile "${PROFILE}" iam put-role-policy \ + --role-name "${ROLE_NAME}" \ + --policy-name "${POLICY_NAME}" \ + --policy-document "${PERMS_POLICY}" + +# Drift check: this script only reconciles the single inline policy above, so +# any OTHER inline or attached policy on the role is an unreviewed grant that +# diverges the live role from these artifacts. Warn-only (never auto-delete) is +# deliberate: this is a shared prod-tenant account and an extra policy might be +# a legitimate out-of-band grant an operator must adjudicate, not blindly strip. +echo "==> Checking for unexpected policies on the role..." +EXTRA_INLINE="$(aws --profile "${PROFILE}" iam list-role-policies \ + --role-name "${ROLE_NAME}" --query "PolicyNames[?@!='${POLICY_NAME}']" --output text)" +EXTRA_ATTACHED="$(aws --profile "${PROFILE}" iam list-attached-role-policies \ + --role-name "${ROLE_NAME}" --query 'AttachedPolicies[].PolicyName' --output text)" +if [[ -n "${EXTRA_INLINE}" || -n "${EXTRA_ATTACHED}" ]]; then + echo " WARNING: unreviewed policies present on ${ROLE_NAME} (not in these artifacts):" >&2 + [[ -n "${EXTRA_INLINE}" ]] && echo " inline: ${EXTRA_INLINE}" >&2 + [[ -n "${EXTRA_ATTACHED}" ]] && echo " attached: ${EXTRA_ATTACHED}" >&2 + echo " Review and remove them (delete-role-policy / detach-role-policy) if unexpected." >&2 +else + echo " OK - only ${POLICY_NAME} present." +fi + +ROLE_ARN="$(aws --profile "${PROFILE}" iam get-role --role-name "${ROLE_NAME}" \ + --query Role.Arn --output text)" + +echo "==> Done. Deploy role ready:" +echo " ${ROLE_ARN}" +echo " Set the GitHub repo secret AWS_DEPLOY_ROLE_ARN to this ARN (Phase 1 merge sequencing:" +echo " flip the secret only when Phase 0 is complete and the migration PR is about to merge)." diff --git a/infra/deploy-role/permissions-policy.json b/infra/deploy-role/permissions-policy.json new file mode 100644 index 0000000..e29ba7e --- /dev/null +++ b/infra/deploy-role/permissions-policy.json @@ -0,0 +1,40 @@ +{ + "Version": "2012-10-17", + "Statement": [ + { + "Sid": "AssumeCdkBootstrapRoles", + "Effect": "Allow", + "Action": "sts:AssumeRole", + "Resource": [ + "arn:aws:iam::011934824531:role/cdk-hnb659fds-deploy-role-011934824531-us-east-1", + "arn:aws:iam::011934824531:role/cdk-hnb659fds-file-publishing-role-011934824531-us-east-1", + "arn:aws:iam::011934824531:role/cdk-hnb659fds-lookup-role-011934824531-us-east-1", + "arn:aws:iam::011934824531:role/cdk-hnb659fds-image-publishing-role-011934824531-us-east-1" + ] + }, + { + "Sid": "CdkDeployHealthCheck", + "Effect": "Allow", + "Action": "cloudformation:DescribeStacks", + "Resource": [ + "arn:aws:cloudformation:us-east-1:011934824531:stack/po-ingest/*", + "arn:aws:cloudformation:us-east-1:011934824531:stack/WorkorderIngestStack/*", + "arn:aws:cloudformation:us-east-1:011934824531:stack/CDKToolkit/*" + ], + "Condition": { + "StringEquals": { + "aws:RequestedRegion": "us-east-1" + } + } + }, + { + "Sid": "SmokeInvokeLambda", + "Effect": "Allow", + "Action": "lambda:InvokeFunction", + "Resource": [ + "arn:aws:lambda:us-east-1:011934824531:function:po-email-processor", + "arn:aws:lambda:us-east-1:011934824531:function:workorder-email-processor" + ] + } + ] +} diff --git a/infra/deploy-role/trust-policy.json b/infra/deploy-role/trust-policy.json new file mode 100644 index 0000000..e67a290 --- /dev/null +++ b/infra/deploy-role/trust-policy.json @@ -0,0 +1,18 @@ +{ + "Version": "2012-10-17", + "Statement": [ + { + "Effect": "Allow", + "Principal": { + "Federated": "arn:aws:iam::011934824531:oidc-provider/token.actions.githubusercontent.com" + }, + "Action": "sts:AssumeRoleWithWebIdentity", + "Condition": { + "StringEquals": { + "token.actions.githubusercontent.com:aud": "sts.amazonaws.com", + "token.actions.githubusercontent.com:sub": "repo:Sea-Haven-Industries/procurement-ingest:ref:refs/heads/main" + } + } + } + ] +} diff --git a/scripts/migrate_tables.py b/scripts/migrate_tables.py new file mode 100644 index 0000000..5f912b3 --- /dev/null +++ b/scripts/migrate_tables.py @@ -0,0 +1,416 @@ +""" +One-off cross-account DynamoDB backfill for the mgmt -> seahaven-prod +migration (Phase 4 of the migration plan). Scans each table in the source +(mgmt) account and writes into the same-named, already-CDK-created table in +the destination (prod) account. + +Dry-run is the DEFAULT in every mode; nothing is written or deleted unless +--execute is passed. + +Per-table modes (encode the plan's dual-delivery-safe semantics): + + WorkOrders, purchase-orders overwrite Put. Safe because during dual + delivery the mgmt item is a strict SUPERSET of + the prod item (both pipelines apply SET-only + merges to the same inbound mail). purchase-orders + additionally PRESERVES a dest-Cancelled PO: a + non-Cancelled source item never overwrites a + prod row already Cancelled (sticky-cancel guard). + WorkOrderComments cutoff-filtered Put: only rows with + attribute_not_exists(ingested_at) OR + ingested_at < --cutoff (T_activate). The same + email yields DIFFERENT comment_ids per account + (S3-object-key hash suffix), so copying + dual-window rows would create duplicate + history. --cutoff is REQUIRED for this table on + BOTH copy and verify. + verified-sites, truncate-and-load: delete ALL destination rows + pending-site-review first, then copy the source scan verbatim. + Required because the purchase-orders backfill + replays the prod site-extractor stream, which + creates junk `pending` rows and non-idempotent + poCount ADDs; mgmt values are ground truth. + Take an on-demand backup of the dest table and + wait for it to be AVAILABLE before running + (the script refuses without a fresh --backup-arn + of THIS table, on DRY-RUN as well as --execute: + the rehearsal validates the same preconditions + the real run will). + +Usage: + # copy (dry-run first, then --execute) -- run tables INDIVIDUALLY, in order + python scripts/migrate_tables.py copy --table WorkOrders + python scripts/migrate_tables.py copy --table WorkOrderComments \ + --cutoff 2026-07-24T02:00:00Z --execute + python scripts/migrate_tables.py copy --table verified-sites \ + --backup-arn arn:aws:dynamodb:...:backup/... --execute + # verify (count parity + spot checks; always read-only) + python scripts/migrate_tables.py verify --table purchase-orders + python scripts/migrate_tables.py verify --all --cutoff 2026-07-24T02:00:00Z + +`copy --all` is deliberately rejected: the plan runs each table as a separate +step with a stream-drain wait before the truncate-and-load tables, so a single +sweep would either write in the wrong order or truncate before the drain. + +Profiles: --source-profile (default: "default", the mgmt SSO session) and +--dest-profile (default: "seahaven-prod"). The script hard-verifies both +resolved account ids before doing anything. +""" + +import argparse +import json +import sys +from datetime import datetime, timedelta, timezone + +import boto3 + +SOURCE_ACCOUNT = "328440206208" +DEST_ACCOUNT = "011934824531" +REGION = "us-east-1" + +# A backup passed to a truncate-and-load must be fresh -- the whole point is a +# just-taken restore point, not a week-old one that predates prod-native rows. +BACKUP_MAX_AGE = timedelta(hours=1) + +TABLES = { + "purchase-orders": {"mode": "overwrite", "keys": ["po_number"]}, + "WorkOrders": {"mode": "overwrite", "keys": ["work_order_id"]}, + "WorkOrderComments": { + "mode": "cutoff", + "keys": ["work_order_id", "comment_id"], + "cutoff_attr": "ingested_at", + }, + "verified-sites": {"mode": "truncate_load", "keys": ["siteCode"]}, + "pending-site-review": {"mode": "truncate_load", "keys": ["po_number"]}, +} + + +def session_table(profile, expected_account, table_name): + 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.resource("dynamodb").Table(table_name) + + +def check_key_schema(table, name, keys): + """Abort with a clear message if the configured keys drift from the live + table's key schema. Every DynamoDB item necessarily carries its table's + key attributes, so once this passes the `item[k]` key projections below + cannot KeyError; a mismatch here is the only way they could. Set (not + ordered) comparison is deliberate: keys are only ever used to build dict + Key= projections, which are order-insensitive.""" + live = {k["AttributeName"] for k in table.key_schema} + if live != set(keys): + sys.exit( + f"ERROR: {name} live key schema {sorted(live)} != configured " + f"{sorted(keys)}; refusing to build keys from a stale TABLES entry." + ) + + +def scan_items(table, filter_kwargs=None): + kwargs = dict(filter_kwargs or {}) + while True: + page = table.scan(**kwargs) + yield from page.get("Items", []) + lek = page.get("LastEvaluatedKey") + if not lek: + return + kwargs["ExclusiveStartKey"] = lek + + +def scan_count(table, filter_kwargs=None): + kwargs = dict(filter_kwargs or {}, Select="COUNT") + total = 0 + while True: + page = table.scan(**kwargs) + total += page["Count"] + lek = page.get("LastEvaluatedKey") + if not lek: + return total + kwargs["ExclusiveStartKey"] = lek + + +def normalize_cutoff(value): + """Parse an operator-supplied cutoff and re-emit the canonical form that + ``ingested_at`` is lexicographically comparable against. + + ``ingested_at`` is ``datetime.now(timezone.utc).isoformat()`` -> a + zero-padded ``...+00:00`` string. A raw operator string ("2026-7-24...", + a space for the 'T', a fractional part) compares WRONG under DynamoDB's + byte-wise ``<`` and would silently copy dual-window rows or drop history. + So parse strictly, require an aware UTC instant, and re-emit second + precision ``+00:00`` with no fraction. A row stamped exactly at T_activate + then sorts >= the cutoff and is excluded -- the safe direction. + """ + raw = value[:-1] + "+00:00" if value.endswith("Z") else value + try: + dt = datetime.fromisoformat(raw) + except ValueError: + sys.exit( + f"ERROR: --cutoff {value!r} is not ISO-8601 (e.g. 2026-07-24T02:00:00Z)." + ) + if dt.tzinfo is None or dt.utcoffset() != timedelta(0): + sys.exit(f"ERROR: --cutoff {value!r} must be UTC (trailing 'Z' or '+00:00').") + return dt.replace(microsecond=0).isoformat() + + +def comments_cutoff_filter(cutoff): + return { + "FilterExpression": ("attribute_not_exists(#ia) OR #ia < :cutoff"), + "ExpressionAttributeNames": {"#ia": "ingested_at"}, + "ExpressionAttributeValues": {":cutoff": cutoff}, + } + + +def _dest_cancelled_po_numbers(dst): + """po_numbers currently Cancelled in the destination purchase-orders table. + + Read once before an overwrite copy so a stale non-Cancelled source item + can never revive a PO the prod pipeline already cancelled (sticky-cancel). + """ + flt = { + "FilterExpression": "po_status = :c", + "ExpressionAttributeValues": {":c": "Cancelled"}, + "ProjectionExpression": "po_number", + } + return {item["po_number"] for item in scan_items(dst, flt)} + + +def _validate_truncate_backup(name, dst, args): + """Fail closed unless --backup-arn is an AVAILABLE, fresh backup of THIS + destination table incarnation (status + name + recency + TableId).""" + if not args.backup_arn: + sys.exit( + f"ERROR: {name} is truncate-and-load; pass --backup-arn of an " + "AVAILABLE on-demand backup of the DESTINATION table taken " + "just now (aws dynamodb create-backup / describe-backup)." + ) + desc = dst.meta.client.describe_backup(BackupArn=args.backup_arn) + details = desc["BackupDescription"]["BackupDetails"] + src_table = desc["BackupDescription"]["SourceTableDetails"] + if details["BackupStatus"] != "AVAILABLE" or src_table["TableName"] != name: + sys.exit( + f"ERROR: backup {args.backup_arn} is {details['BackupStatus']} for " + f"table {src_table['TableName']!r}; need an AVAILABLE backup of " + f"{name!r}. Aborting before truncate." + ) + live_id = dst.meta.client.describe_table(TableName=name)["Table"]["TableId"] + if src_table.get("TableId") != live_id: + sys.exit( + f"ERROR: backup {args.backup_arn} is of a DIFFERENT incarnation of " + f"{name!r} (TableId {src_table.get('TableId')} != live {live_id}); " + "a RETAIN+recreate cycle produces a same-named table. Aborting." + ) + age = datetime.now(timezone.utc) - details["BackupCreationDateTime"] + if age > BACKUP_MAX_AGE: + sys.exit( + f"ERROR: backup {args.backup_arn} is {age} old (> {BACKUP_MAX_AGE}); " + "take a fresh backup immediately before truncating. Aborting." + ) + + +def _truncate_destination(name, dst, args, keys): + # Validate in dry-run too: a stale/missing --backup-arn should surface on + # the rehearsal run, not first appear when the operator adds --execute. + _validate_truncate_backup(name, dst, args) + dest_keys = [{k: item[k] for k in keys} for item in scan_items(dst)] + print(f"[{name}] truncate: {len(dest_keys)} destination rows to delete") + if args.execute: + with dst.batch_writer() as batch: + for key in dest_keys: + batch.delete_item(Key=key) + print(f"[{name}] truncate: done") + else: + print(f"[{name}] DRY-RUN: no deletes performed") + + +def copy_table(name, cfg, args): + src = session_table(args.source_profile, SOURCE_ACCOUNT, name) + dst = session_table(args.dest_profile, DEST_ACCOUNT, name) + check_key_schema(src, name, cfg["keys"]) + check_key_schema(dst, name, cfg["keys"]) + + filter_kwargs = None + if cfg["mode"] == "cutoff": + # --cutoff presence is enforced up front in main(); normalize here. + filter_kwargs = comments_cutoff_filter(normalize_cutoff(args.cutoff)) + + if cfg["mode"] == "truncate_load": + _truncate_destination(name, dst, args, cfg["keys"]) + + # Sticky-cancel guard: never let a stale non-Cancelled source item revive a + # PO already Cancelled in prod (the raw batch Put bypasses the handler's + # ConditionExpression that normally enforces this). + dest_cancelled = ( + _dest_cancelled_po_numbers(dst) if name == "purchase-orders" else set() + ) + + copied = preserved = 0 + if args.execute: + with dst.batch_writer() as batch: + for item in scan_items(src, filter_kwargs): + if ( + item.get("po_number") in dest_cancelled + and item.get("po_status") != "Cancelled" + ): + preserved += 1 + continue + batch.put_item(Item=item) + copied += 1 + else: + for item in scan_items(src, filter_kwargs): + if ( + item.get("po_number") in dest_cancelled + and item.get("po_status") != "Cancelled" + ): + preserved += 1 + continue + copied += 1 + verb = "copied" if args.execute else "would copy (DRY-RUN)" + print(f"[{name}] {verb} {copied} items {SOURCE_ACCOUNT} -> {DEST_ACCOUNT}") + if preserved: + print(f"[{name}] preserved {preserved} dest-Cancelled PO(s) (not overwritten)") + + +def verify_table(name, cfg, args): + src = session_table(args.source_profile, SOURCE_ACCOUNT, name) + dst = session_table(args.dest_profile, DEST_ACCOUNT, name) + check_key_schema(src, name, cfg["keys"]) + check_key_schema(dst, name, cfg["keys"]) + + src_filter = None + if cfg["mode"] == "cutoff": + # main() guarantees --cutoff is present when a cutoff table is in scope. + src_filter = comments_cutoff_filter(normalize_cutoff(args.cutoff)) + + src_count = scan_count(src, src_filter) + dst_count = scan_count(dst) + # Destination may legitimately EXCEED source-side-of-cutoff (prod-native + # dual-delivery rows) for cutoff/overwrite tables; exact parity is expected + # only for truncate_load immediately after the load. + relation = { + "overwrite": dst_count >= src_count, + "cutoff": dst_count >= src_count, + "truncate_load": dst_count == src_count, + }[cfg["mode"]] + status = "OK" if relation else "MISMATCH" + print(f"[{name}] source={src_count} dest={dst_count} -> {status}") + + ok = relation + # Spot-check every mode. Key-existence lookups are sound for ALL modes + # because the copy batch-Puts source items VERBATIM (source comment_id + # included) -- per-account comment_id divergence only affects rows each + # pipeline ingests itself, and those dual-window rows are excluded here by + # the same cutoff filter the copy used. For truncate_load keys are likewise + # copied verbatim. + if args.sample: + sampled = mismatched = 0 + for item in scan_items(src, src_filter): + if sampled >= args.sample: + break + sampled += 1 + key = {k: item[k] for k in cfg["keys"]} + got = dst.get_item(Key=key).get("Item") + if got is None: + mismatched += 1 + print(f"[{name}] MISSING in dest: {json.dumps(key, default=str)}") + print(f"[{name}] spot-check: {sampled} sampled, {mismatched} missing") + ok = ok and mismatched == 0 + + cancelled_check_limit = 25 + if name == "purchase-orders": + # Propagation check: a PO Cancelled in mgmt must be Cancelled in prod. + # (Revival in the OTHER direction -- a prod-Cancelled PO overwritten by + # a non-Cancelled mgmt item -- is prevented copy-side by the sticky- + # cancel guard, not detectable here after the fact.) + cancelled_filter = { + "FilterExpression": "po_status = :c", + "ExpressionAttributeValues": {":c": "Cancelled"}, + } + checked = wrong = 0 + for item in scan_items(src, cancelled_filter): + if checked >= cancelled_check_limit: + break + checked += 1 + got = dst.get_item(Key={"po_number": item["po_number"]}).get("Item") + if not got or got.get("po_status") != "Cancelled": + wrong += 1 + print(f"[{name}] CANCELLED-DRIFT: {item['po_number']}") + print(f"[{name}] cancelled spot-check: {checked} checked, {wrong} drifted") + ok = ok and wrong == 0 + + return ok + + +def _require_cutoff_if_needed(names, args): + cutoff_tables = [n for n in names if TABLES[n]["mode"] == "cutoff"] + if cutoff_tables and not args.cutoff: + sys.exit( + f"ERROR: {', '.join(cutoff_tables)} is cutoff-mode; --cutoff " + "(T_activate) is required so dual-window rows are excluded." + ) + + +def main(): + ap = argparse.ArgumentParser(description=__doc__) + sub = ap.add_subparsers(dest="cmd", required=True) + for cmd in ("copy", "verify"): + p = sub.add_parser(cmd) + p.add_argument("--table", choices=sorted(TABLES)) + p.add_argument( + "--all", + action="store_true", + help="all five tables (verify only; copy runs tables individually)", + ) + p.add_argument("--source-profile", default="default") + p.add_argument("--dest-profile", default="seahaven-prod") + p.add_argument( + "--cutoff", + help="T_activate ISO-8601 UTC; required for WorkOrderComments", + ) + if cmd == "copy": + p.add_argument("--execute", action="store_true") + p.add_argument( + "--backup-arn", + help="required for truncate-and-load tables (validated on " + "dry-run as well as --execute)", + ) + else: + p.add_argument( + "--sample", + type=int, + default=10, + help="per-table spot-check sample size (0 disables)", + ) + + args = ap.parse_args() + if bool(args.table) == bool(args.all): + sys.exit("ERROR: pass exactly one of --table or --all.") + + if args.cmd == "copy" and args.all: + # Validation precedes any side effect: --all would copy in the wrong + # order and truncate before the stream drain. Run tables individually. + sys.exit( + "ERROR: `copy --all` is not supported. Run each table individually " + "in plan order (WorkOrders, WorkOrderComments, purchase-orders, " + "wait for the site-extractor stream to drain, then verified-sites " + "and pending-site-review)." + ) + + names = sorted(TABLES) if args.all else [args.table] + _require_cutoff_if_needed(names, args) + + if args.cmd == "copy": + copy_table(args.table, TABLES[args.table], args) + else: + results = [verify_table(n, TABLES[n], args) for n in names] + if not all(results): + sys.exit(1) + + +if __name__ == "__main__": + main() diff --git a/tests/test_reprocess_contract.py b/tests/test_reprocess_contract.py index a899966..81de68a 100644 --- a/tests/test_reprocess_contract.py +++ b/tests/test_reprocess_contract.py @@ -44,7 +44,9 @@ _spec.loader.exec_module(reprocess) def test_build_s3_event_shape_and_raw_key(): - bucket = "po-ingest-emails-328440206208" + # Account-agnostic fixture (the suffix is never parsed); the real bucket + # name is account-derived at deploy time. + bucket = "po-ingest-emails-000000000000" # Deliberately contains a SPACE, a literal '+', and a literal '%41'. # A real S3 event notification would deliver this key encoded as # "inbound/2026/AB+12%2B34+%2541.eml"