mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 11:53:13 +00:00
417 lines
17 KiB
Python
417 lines
17 KiB
Python
|
|
"""
|
||
|
|
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()
|