Files
jorgecuadros-platform/migration/run_all.py
T
rmancinasandClaude Opus 5 ed19f51a52
Build and Push Images / Deploy to galactus (push) Canceled after 0s
Build and Push Images / Build jorgecuadros-web (push) Canceled after 1m26s
Build and Push Images / Build jorgecuadros-api (push) Canceled after 1m27s
fix(ops): re-stage before an additive sync
SYNC ran `run_all.py --sync` without `--stage`, so it depended on staged
Parquet under migration/output. That directory is part of the image, not a
volume, so any redeploy wiped it and the job died on the first transform:

    FileNotFoundError: '/repo/migration/output/stg_utilities/datgral.parquet'

Re-staging is also what makes the job's own label true — without it a sync
would replay whatever upload staged last, not the files currently in the
ingest folder.

Staging now counts as a numbered step when it runs, so the Operaciones
progress bar moves during the slowest phase instead of sitting empty.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 22:20:30 -07:00

130 lines
5.5 KiB
Python

"""
Run the full data migration against one environment, in dependency order.
Every step is idempotent (truncate + rebuild), so this is safe to re-run. The
target DB is chosen with --env (reads deploy/.env.<env>); the same staged
Parquet feeds every environment.
Prerequisites (once per environment, NOT done here):
1. DB stack deployed (deploy/jorgecuadros-db.stack.yml) and deploy/.env.<env> written.
2. Prisma schema pushed to it:
DATABASE_URL="<that env's url>" \
npx prisma@5 db push --schema=packages/database/prisma/schema.prisma
Then:
./.venv/bin/python run_all.py --env dev # data only (staging already present)
./.venv/bin/python run_all.py --env prod --stage # re-extract from Access first, then load
--sync swaps the truncate+rebuild steps for the additive upsert ones. It reads
the same staged Parquet, so it needs --stage too unless a previous run left
migration/output populated on this machine — which is never true in a
container, where that directory is part of the image and dies with it.
Reproducing dev -> prod is exactly `--env prod` (plus --stage if the staged
Parquet isn't present on the machine running it).
Note: the blob_extract step re-reads the original Access files directly (the
blobs are not in the staged Parquet), so the machine running this needs
SOURCE_ROOT + mdbtools + MinIO credentials even without --stage.
"""
from __future__ import annotations
import argparse
import subprocess
import sys
from pathlib import Path
HERE = Path(__file__).parent
PY = sys.executable # the venv python running this orchestrator
# Dependency order. Every step truncates what it owns, so anything downstream
# of a truncated table has to be rebuilt in the same pass — blob_extract is in
# this list because transform_properties and transform_policies truncate
# service_documents / policy_documents, which would otherwise leave the
# uploaded MinIO objects with no rows pointing at them.
STEPS = [
"transform_customers.py",
"transform_properties.py",
# Statement-OCR match fields. transform_properties.py now produces these
# directly, so on a full rebuild this is a no-op that re-asserts they are
# there; on a database predating the OCR module it is what fills them in.
# Must follow transform_properties.py, which truncates both tables it
# touches.
"backfill_statement_match_fields.py",
"transform_policies.py",
"transform_transactions.py",
"prune_empty_customers.py",
# Seeds the Scotiabank chequera that every SCOTHIA movement is booked into;
# transform_bank.py fails fast without it.
"backfill_bank_accounts.py",
"transform_bank.py",
"blob_extract.py",
]
SYNC_STEPS = [
"transform_customers.py",
"transform_properties.py",
# Statement-OCR match fields. transform_properties.py now produces these
# directly, so on a full rebuild this is a no-op that re-asserts they are
# there; on a database predating the OCR module it is what fills them in.
# Must follow transform_properties.py, which truncates both tables it
# touches.
"backfill_statement_match_fields.py",
"transform_policies.py",
"transform_transactions.py",
# Manual-safe prune: drops legacy-owned empties that the customer upsert
# re-creates from Parquet, but leaves manually-added customers alone.
"prune_empty_customers.py",
# Seeds the Scotiabank chequera that every SCOTHIA movement is booked into;
# transform_bank.py fails fast without it.
"backfill_bank_accounts.py",
"transform_bank.py",
]
def run(cmd: list[str], step: int | None = None, total: int | None = None) -> None:
# The "[paso i/N] name" marker is a contract with the Operaciones screen,
# which parses the last one to show progress. Emitting it here rather than
# letting the UI count STEPS itself keeps the two from drifting when a step
# is added — the number of steps is only ever stated in this file.
if step is not None and total is not None:
print(f"[paso {step}/{total}] {Path(cmd[1]).name}", flush=True)
print("+ " + " ".join(cmd), flush=True)
r = subprocess.run(cmd)
if r.returncode:
sys.exit(r.returncode)
def main() -> None:
ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--env", default="dev", help="target environment (reads deploy/.env.<env>)")
ap.add_argument("--stage", action="store_true",
help="re-run the raw staging load first (needs the Access files + mdbtools)")
ap.add_argument("--sync", action="store_true",
help="upsert legacy rows and archive removed legacy rows; preserve manual rows")
args = ap.parse_args()
steps = SYNC_STEPS if args.sync else STEPS
# Staging counts as a step when it runs: it is the slowest part of the pass
# (mdbtools re-reads every Access file), so leaving it outside the numbering
# would park the Operaciones progress bar at "nothing yet" for minutes.
total = len(steps) + (1 if args.stage else 0)
offset = 1 if args.stage else 0
if args.stage:
run([PY, str(HERE / "load_staging.py"), "--output-dir", str(HERE / "output")],
step=1, total=total)
for i, step in enumerate(steps, start=1 + offset):
cmd = [PY, str(HERE / step), "--env", args.env]
if args.sync:
cmd.append("--sync")
run(cmd, step=i, total=total)
print(f"\n✓ migration complete for env={args.env}")
if __name__ == "__main__":
main()