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