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