Files
rmancinasandClaude Opus 5 898cf48c80 fix(migration): re-import died on the last step because blob_extract required deploy/.env.prod
Every transform resolves its target through dbenv.database_url(), which lets a
DATABASE_URL in the process environment win — that is how the API container
drives a re-import against its own database with no deploy/ directory present.
blob_extract.py was the one step that bypassed it and called load_env()
directly for the MinIO credentials, so the "Operaciones" re-import loaded all
the data and then exited 1 on:

  missing /repo/deploy/.env.prod — deploy the 'prod' DB stack and write its
  .env first

Give the S3 settings the same resolution as the DB URL: load_env() now returns
{} for an absent file, and setting()/require() layer the process environment on
top of it. blob_extract reads S3_ENDPOINT / S3_BUCKET and accepts either
S3_ACCESS_KEY/S3_SECRET_KEY or MINIO_ROOT_USER/MINIO_ROOT_PASSWORD, matching
the fallback order in storage.service.ts and the vars the api service already
sets in deploy/galactus/jorgecuadros-app.compose.yml. A genuinely missing
setting still fails fast, now naming the variable.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-01 02:21:38 -07:00

217 lines
8.9 KiB
Python

"""
Migration plan step 4: extract LONGBINARY document blobs to object storage.
The Access LONGBINARY columns hold scanned utility bills / IDs / policy docs
wrapped in an Access OLE Object container (a "\\x15\\x1c...Pres..." header +
optional DIB preview, then the real embedded file). Staging used
`mdb-export -b strip` (blobs dropped); this re-reads each table with
`-b hex`, carves the embedded file out of the OLE wrapper by locating its
magic bytes, uploads it to MinIO (S3), and writes a *_documents row pointing
at it (MySQL keeps only the pointer + metadata, per the plan).
Row alignment: `mdb-export` order is deterministic and identical to the order
load_staging used, so a row's position == its staged `_row_num`. Policy docs
resolve to a policy by (legacySourceTable, legacyId=row position); property
docs resolve by DATMEX numer_id -> propertyId.
Idempotent: truncates the *_documents tables for the selected models and
re-uploads under deterministic keys (overwrite). Use --limit N for a small
test pass, --tables to restrict.
Run: ./.venv/bin/python blob_extract.py --env dev [--limit N] [--tables datmex,mult]
"""
from __future__ import annotations
import argparse
import csv
import subprocess
import sys
import uuid
from pathlib import Path
import boto3
from botocore.config import Config
from dbenv import connect, require, setting
from extract import sanitize_column_name as san
csv.field_size_limit(300_000_000)
# Same source folder as the rest of the pipeline (config.SOURCE_ROOT honours
# INGEST_DIR — the web "Operaciones" ingest volume).
from config import SOURCE_ROOT
# (key, access_file, access_table, staged_table_name, blob_cols, model)
SOURCES = [
# DATMEX's real scanned bills live in the ILUZ/IAGUA/IPREDIAL/ITEL invoice-
# image columns (per-service bill scans), not doc_1/doc_2 (which are empty).
# In practice only a handful are populated — the .accdb is mostly bloat.
dict(key="datmex", file="UTILITIES.accdb", table="DATMEX", staged="DATMEX",
blobs=["iluz", "iagua", "ipredial", "itel", "doc_1", "doc_2"], model="service",
doctypes={"iluz": "ELECTRIC_BILL", "iagua": "WATER_BILL",
"ipredial": "PROPERTY_TAX_BILL", "itel": "PHONE_BILL"}),
dict(key="mult", file="SEGUROS 16_be.mdb", table="MULT", staged="mult",
blobs=["foto1", "docs_1", "docs_2"], model="policy"),
dict(key="autos_ampl", file="SEGUROS 16_be.mdb", table="TABLA AUTOS AMPL",
staged="tabla_autos_ampl", blobs=["foto1", "docs_1", "docs_2"], model="policy"),
]
# magic -> (ext, content-type). Order = priority when several appear.
MAGICS = [
(b"\xff\xd8\xff", "jpg", "image/jpeg"),
(b"\x89PNG\r\n\x1a\n", "png", "image/png"),
(b"%PDF", "pdf", "application/pdf"),
(b"GIF8", "gif", "image/gif"),
(b"II*\x00", "tif", "image/tiff"),
(b"MM\x00*", "tif", "image/tiff"),
]
def carve(b: bytes):
"""Locate the embedded file inside the OLE wrapper and return
(bytes, ext, content_type) or None if no known type is present."""
best = None
for sig, ext, ct in MAGICS:
i = b.find(sig)
if i >= 0 and (best is None or i < best[0]):
best = (i, ext, ct)
if best is None:
return None
i, ext, ct = best
data = b[i:]
# trim trailing OLE junk after the real end marker where we know it
if ext == "jpg":
e = data.rfind(b"\xff\xd9")
if e >= 0:
data = data[: e + 2]
elif ext == "png":
e = data.rfind(b"IEND\xaeB`\x82")
if e >= 0:
data = data[: e + 8]
return data, ext, ct
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--env", default="dev")
ap.add_argument("--limit", type=int, default=0, help="max rows per table (0 = all); test with a small N")
ap.add_argument("--tables", default="", help="comma list of source keys to run (default all)")
args = ap.parse_args()
only = set(x.strip() for x in args.tables.split(",") if x.strip())
sources = [s for s in SOURCES if not only or s["key"] in only]
# Process environment first, deploy/.env.<env> second, with the same
# credential aliases the API uses — the "Operaciones" re-import runs this
# inside the API container, which has S3_ENDPOINT / MINIO_ROOT_* injected
# and no deploy/ directory at all.
s3 = boto3.client(
"s3", endpoint_url=require(args.env, "S3_ENDPOINT"),
aws_access_key_id=require(args.env, "S3_ACCESS_KEY", "MINIO_ROOT_USER"),
aws_secret_access_key=require(args.env, "S3_SECRET_KEY", "MINIO_ROOT_PASSWORD"),
config=Config(signature_version="s3v4"), region_name="us-east-1")
bucket = setting(args.env, "S3_BUCKET") or "jorgecuadros-documents"
conn = connect(args.env)
cur = conn.cursor()
# parent lookups
cur.execute("SELECT legacyId, id FROM properties WHERE legacySourceTable='DATMEX'")
prop_by_numer = {}
for lid, pid in cur.fetchall():
prop_by_numer.setdefault(lid, pid) # first property per numer_id
cur.execute("SELECT legacySourceTable, legacyId, id FROM policies")
pol_by_row = {(t, l): i for t, l, i in cur.fetchall()}
# fresh rebuild of the doc tables we're loading (unless a limited test pass)
if not args.limit:
cur.execute("SET FOREIGN_KEY_CHECKS=0")
if any(s["model"] == "service" for s in sources):
cur.execute("TRUNCATE TABLE service_documents")
if any(s["model"] == "policy" for s in sources):
cur.execute("TRUNCATE TABLE policy_documents")
cur.execute("SET FOREIGN_KEY_CHECKS=1")
conn.commit()
svc_rows, pol_rows = [], []
stats = {}
for src in sources:
path = SOURCE_ROOT / src["file"]
want = {c: None for c in src["blobs"]}
uploaded = skipped_noparent = no_magic = empty = 0
p = subprocess.Popen(["mdb-export", "-b", "hex", str(path), src["table"]],
stdout=subprocess.PIPE, text=True, encoding="utf-8",
errors="replace", bufsize=1)
rdr = csv.reader(p.stdout)
hdr = next(rdr)
sh = [san(c) for c in hdr]
idx = {c: sh.index(c) for c in src["blobs"] if c in sh}
numer_idx = sh.index("numer_id") if "numer_id" in sh else None
for ri, row in enumerate(rdr):
if args.limit and ri >= args.limit:
break
# resolve parent
if src["model"] == "service":
numer = (row[numer_idx].strip() if numer_idx is not None and numer_idx < len(row) else "")
if numer.endswith(".0"):
numer = numer[:-2]
parent = prop_by_numer.get(numer)
else:
parent = pol_by_row.get((src["staged"], str(ri)))
for col, ci in idx.items():
h = row[ci].strip() if ci < len(row) else ""
if len(h) < 16:
empty += 1
continue
if not parent:
skipped_noparent += 1
continue
try:
raw = bytes.fromhex(h)
except ValueError:
continue
out = carve(raw)
if out is None:
no_magic += 1
continue
data, ext, ct = out
prefix = "service" if src["model"] == "service" else "policy"
key = f"{prefix}/{parent}/{src['staged']}_{ri}_{col}.{ext}"
s3.put_object(Bucket=bucket, Key=key, Body=data, ContentType=ct)
dtype = src.get("doctypes", {}).get(col, col.upper())
if src["model"] == "service":
svc_rows.append((str(uuid.uuid4()), parent, dtype, key))
else:
pol_rows.append((str(uuid.uuid4()), parent, dtype, key, col))
uploaded += 1
p.stdout.close(); p.wait()
stats[src["key"]] = dict(uploaded=uploaded, no_parent=skipped_noparent,
no_magic=no_magic, empty_cells=empty)
print(f" [{src['key']}] uploaded={uploaded} no_parent={skipped_noparent} "
f"no_magic={no_magic}")
if svc_rows:
cur.executemany("INSERT INTO service_documents (id,propertyId,documentType,storageKey) "
"VALUES (%s,%s,%s,%s)", svc_rows)
if pol_rows:
cur.executemany("INSERT INTO policy_documents (id,policyId,documentType,storageKey,originalColumn) "
"VALUES (%s,%s,%s,%s,%s)", pol_rows)
conn.commit()
cur.execute("SELECT COUNT(*) FROM service_documents"); ns = cur.fetchone()[0]
cur.execute("SELECT COUNT(*) FROM policy_documents"); npd = cur.fetchone()[0]
print("=== Blob extraction complete ===")
for k, v in stats.items():
print(f" {k}: {v}")
print(f" -> service_documents rows total: {ns}")
print(f" -> policy_documents rows total : {npd}")
conn.close()
if __name__ == "__main__":
main()