#!/usr/bin/env python3 """ scripts/migrate_photos_to_r2.py — Verified per-tenant photo sync to R2 ======================================================================== Multi-tenant port of the single-tenant script of the same name (ST §22 Phase 3). Copies each tenant's photo files from the shared local uploads tree to the R2 bucket **under that tenant's object prefix**, and VERIFIES each one (MD5 + size). It is: - COPY-ONLY : never deletes or modifies local files, never touches any DB. - IDEMPOTENT : re-running skips objects already present and verified. - RESUMABLE : safe to interrupt and re-run; picks up where it left off. - GATED : exits non-zero if ANY file fails verification, so you never flip STORAGE_BACKEND=s3 on an incomplete/corrupt copy. Key mapping — the whole point of this script -------------------------------------------- The DB stores an unprefixed key and always will: DB value : uploads/issue_photos/ab12….jpg R2 object : t3/uploads/issue_photos/ab12….jpg The ``t/`` prefix is applied by ``S3Backend._object_key()`` at runtime (app/utils/storage.py). This script must therefore write to the SAME prefixed key, or nothing will resolve after cutover. Because the local backend does NOT prefix, one shared directory holds every tenant's files and a file on disk carries no ownership marker. Ownership is resolved exactly as in scripts/audit_photos.py: by reading each tenant's own database and collecting the keys it references. Consequences, both intended: * Files referenced by NO tenant (orphans) are NOT uploaded — no tenant's prefix could legitimately claim them, and nothing reads them. * A key referenced by two tenants is copied under BOTH prefixes. Run audit_photos.py first; it flags that case. R2 credentials come from the environment / .env: R2_ENDPOINT_URL, R2_ACCESS_KEY_ID, R2_SECRET_ACCESS_KEY, R2_BUCKET. Control-plane access needs CONTROL_DATABASE_URL and CONTROL_FERNET_KEY. Run this while the app is still serving on STORAGE_BACKEND=local — it only writes to R2. Why put_object (not upload_file): a single-part PUT makes the R2 ETag equal the object's MD5, so verification is a direct hash comparison. Photos are <= 50 MB (MAX_CONTENT_LENGTH), well within a single PUT. Usage ----- cd /home/jqc/lt_janitorial_quality_control source venv/bin/activate pip install boto3 # if not already python scripts/migrate_photos_to_r2.py --dry-run # plan only python scripts/migrate_photos_to_r2.py # all tenants python scripts/migrate_photos_to_r2.py --tenant lts # one tenant python scripts/migrate_photos_to_r2.py # again at cutover (delta) python scripts/migrate_photos_to_r2.py --force # re-upload everything python scripts/migrate_photos_to_r2.py --log /home/jqc/r2_sync.json Exit code: 0 only when every referenced file that exists on disk is verified in R2 under its tenant prefix, with zero mismatches and zero errors. Non-zero otherwise. References missing on disk are reported and do NOT fail the gate — they were already broken before this move (audit_photos.py is the baseline). """ import os import sys import json import hashlib import argparse from datetime import datetime _REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), '..')) sys.path.insert(0, _REPO_ROOT) from dotenv import load_dotenv # noqa: E402 load_dotenv(os.path.join(_REPO_ROOT, '.env')) # Reuse the audit script's reference collection so the two can never diverge. from scripts.audit_photos import ( # noqa: E402 AUDITABLE_STATUSES, STATIC_FOLDER, UPLOADS_ROOT, collect_referenced, select_tenants, ) _CONTENT_TYPES = { 'jpg': 'image/jpeg', 'jpeg': 'image/jpeg', 'png': 'image/png', 'gif': 'image/gif', } _CHUNK = 1024 * 1024 # 1 MB def _md5_and_size(path): h = hashlib.md5() size = 0 with open(path, 'rb') as fh: while True: chunk = fh.read(_CHUNK) if not chunk: break h.update(chunk) size += len(chunk) return h.hexdigest(), size def _content_type(key): ext = key.rsplit('.', 1)[-1].lower() if '.' in key else '' return _CONTENT_TYPES.get(ext, 'application/octet-stream') def _remote_state(client, bucket, key): """Return (exists, etag_no_quotes, size) for a key in the bucket.""" try: resp = client.head_object(Bucket=bucket, Key=key) return True, resp['ETag'].strip('"'), resp['ContentLength'] except Exception: return False, None, None def human(n): n = float(n) for unit in ('B', 'KB', 'MB', 'GB', 'TB'): if abs(n) < 1024.0: return f'{n:.1f} {unit}' n /= 1024.0 return f'{n:.1f} PB' def main(): parser = argparse.ArgumentParser( description='Verified per-tenant photo sync to R2 (prefix t/).') parser.add_argument('--tenant', default='all', help="Tenant id or slug, or 'all' (default).") parser.add_argument('--dry-run', action='store_true', help='Report what would happen; upload nothing.') parser.add_argument('--force', action='store_true', help='Re-upload even if the object already verifies.') parser.add_argument('--log', default=None, help='Path to write the JSON result log (default: /tmp/...).') args = parser.parse_args() missing_cfg = [k for k in ('R2_ENDPOINT_URL', 'R2_ACCESS_KEY_ID', 'R2_SECRET_ACCESS_KEY', 'R2_BUCKET') if not os.environ.get(k)] if missing_cfg: print(f'ERROR: missing R2 config: {", ".join(missing_cfg)}') print('Set them in .env before running the sync.') sys.exit(2) try: import boto3 from botocore.config import Config as BotoConfig except ImportError: print('ERROR: boto3 not installed. Run: pip install boto3') sys.exit(2) bucket = os.environ['R2_BUCKET'] client = boto3.client( 's3', endpoint_url = os.environ['R2_ENDPOINT_URL'], aws_access_key_id = os.environ['R2_ACCESS_KEY_ID'], aws_secret_access_key = os.environ['R2_SECRET_ACCESS_KEY'], region_name = 'auto', config = BotoConfig(signature_version='s3v4'), ) if not os.path.isdir(UPLOADS_ROOT): print(f'ERROR: uploads folder not found: {UPLOADS_ROOT}') sys.exit(2) tenants = select_tenants(args.tenant) if not tenants: print(f'No tenants matched --tenant {args.tenant!r} ' f'(statuses: {", ".join(AUDITABLE_STATUSES)}).') sys.exit(2) totals = {'referenced': 0, 'missing_on_disk': 0, 'skipped_verified': 0, 'uploaded': 0, 'verify_failed': 0, 'errors': 0, 'bytes_uploaded': 0} failures = [] per_tenant = [] print('=' * 72) print(f' R2 PHOTO SYNC — MULTI-TENANT{" (DRY RUN)" if args.dry_run else ""}') print(f' bucket : {bucket}') print(f' source : {UPLOADS_ROOT}') print(f' tenants : {len(tenants)}') print('=' * 72) for tid, slug, db_uri in tenants: prefix = f't{tid}/' counts = {'referenced': 0, 'missing_on_disk': 0, 'skipped_verified': 0, 'uploaded': 0, 'verify_failed': 0, 'errors': 0, 'bytes_uploaded': 0} print(f'-- tenant {tid} ({slug}) prefix={prefix}') try: referenced, _by_source = collect_referenced(db_uri) except Exception as exc: print(f' ERROR reading tenant DB: {exc}') totals['errors'] += 1 failures.append({'tenant_id': tid, 'slug': slug, 'stage': 'tenant_db', 'error': str(exc)}) per_tenant.append({'tenant_id': tid, 'slug': slug, 'error': str(exc)}) continue for key in sorted(referenced.keys()): counts['referenced'] += 1 abs_path = os.path.normpath(os.path.join(STATIC_FOLDER, key)) if not os.path.isfile(abs_path): # Pre-existing broken reference — surfaced by audit_photos.py. counts['missing_on_disk'] += 1 continue try: local_md5, local_size = _md5_and_size(abs_path) except OSError as e: counts['errors'] += 1 failures.append({'tenant_id': tid, 'key': key, 'stage': 'read', 'error': str(e)}) print(f' ERROR read {key}: {e}') continue object_key = f'{prefix}{key}' # Resumable skip: already present and matching? if not args.force: exists, r_etag, r_size = _remote_state(client, bucket, object_key) if exists and r_etag == local_md5 and r_size == local_size: counts['skipped_verified'] += 1 continue if args.dry_run: counts['uploaded'] += 1 # would upload print(f' WOULD PUT {object_key} ({human(local_size)})') continue # Upload (single-part PUT so ETag == MD5) try: with open(abs_path, 'rb') as body: client.put_object( Bucket=bucket, Key=object_key, Body=body, ContentType=_content_type(key), ) except Exception as e: counts['errors'] += 1 failures.append({'tenant_id': tid, 'key': object_key, 'stage': 'upload', 'error': str(e)}) print(f' ERROR upload {object_key}: {e}') continue # Verify: re-HEAD and compare ETag(MD5) + size exists, r_etag, r_size = _remote_state(client, bucket, object_key) if exists and r_etag == local_md5 and r_size == local_size: counts['uploaded'] += 1 counts['bytes_uploaded'] += local_size else: counts['verify_failed'] += 1 failures.append({ 'tenant_id': tid, 'key': object_key, 'stage': 'verify', 'local_md5': local_md5, 'local_size': local_size, 'remote_etag': r_etag, 'remote_size': r_size, }) print(f' VERIFY FAIL {object_key} local_md5={local_md5} ' f'remote_etag={r_etag} local_size={local_size} remote_size={r_size}') for k in totals: totals[k] += counts[k] per_tenant.append({'tenant_id': tid, 'slug': slug, 'prefix': prefix, 'counts': counts}) verified_here = counts['skipped_verified'] + (0 if args.dry_run else counts['uploaded']) print(f' referenced={counts["referenced"]} ' f'missing_on_disk={counts["missing_on_disk"]} ' f'verified={verified_here} ' f'uploaded={counts["uploaded"]} ({human(counts["bytes_uploaded"])}) ' f'verify_failed={counts["verify_failed"]} errors={counts["errors"]}') # ── summary ── print('-' * 72) if args.dry_run: action_line = f' would upload : {totals["uploaded"]}' else: action_line = (f' uploaded + verified : {totals["uploaded"]}' f' ({human(totals["bytes_uploaded"])})') print(f' referenced keys (all tenants) : {totals["referenced"]}') print(f' referenced but missing on disk: {totals["missing_on_disk"]} ' f'(pre-existing broken refs)') print(f' already verified (skipped) : {totals["skipped_verified"]}') print(action_line) print(f' verify failures : {totals["verify_failed"]}') print(f' errors : {totals["errors"]}') print('=' * 72) verified_total = totals['skipped_verified'] + (0 if args.dry_run else totals['uploaded']) syncable = totals['referenced'] - totals['missing_on_disk'] safe = (totals['verify_failed'] == 0 and totals['errors'] == 0) if args.dry_run: print(' DRY RUN — nothing uploaded. Re-run without --dry-run to sync.') elif safe: print(f' GATE: {verified_total}/{syncable} referenced files verified in R2 under ' f'their tenant prefix, 0 mismatches, 0 errors.') print(' ✅ SAFE to flip STORAGE_BACKEND=s3. Local files are untouched.') else: print(f' GATE: {totals["verify_failed"]} verify failure(s), {totals["errors"]} error(s).') print(' ❌ NOT safe to cut over. Investigate failures (see log), then re-run.') log_path = args.log or os.path.join( '/tmp', f'jqc_r2_sync_{datetime.utcnow().strftime("%Y%m%d_%H%M%S")}.json') try: with open(log_path, 'w') as fh: json.dump({ 'generated_at': datetime.utcnow().isoformat() + 'Z', 'bucket': bucket, 'dry_run': args.dry_run, 'tenant_selector': args.tenant, 'totals': totals, 'tenants': per_tenant, 'failures': failures, }, fh, indent=2) print(f'\n Log written to: {log_path}') except OSError as e: print(f'\n Could not write log to {log_path}: {e}') sys.exit(0 if safe else 1) if __name__ == '__main__': main()