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).
This commit is contained in:
Adam Moussa 2026-07-23 17:08:47 -04:00 • committed by GitHub
parent 820f86ff2e
commit 073201f633
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 636 additions and 8 deletions

View file

@ -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

View file

@ -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**

View file

@ -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).

View file

@ -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/<queue-name>`.
DLQ URLs are `https://sqs.us-east-1.amazonaws.com/011934824531/<queue-name>`.
(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)

View file

@ -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.

View file

@ -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)."

View file

@ -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"
]
}
]
}

View file

@ -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"
}
}
}
]
}

416
scripts/migrate_tables.py Normal file
View file

@ -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()

View file

@ -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"