#!/usr/bin/env python3 """ move_project.py — Move a music project between vendor accounts in art_relations. Setup: python3 -m venv .venv source .venv/bin/activate pip install -r requirements.txt Usage: python move_project.py --project-id 12345 --destination-vendor-id 99 python move_project.py --project-id 12345 --destination-vendor-id 99 --destination-subaccount-id 7 python move_project.py --project-id 12345 --destination-vendor-id 99 --dry-run Required env vars (or .env file): MYSQL_HOST, MYSQL_USER, MYSQL_PASSWORD GRAPHQL_URL GRAPHQL_ORCHARD_USER_ID GRAPHQL_ORCHARD_IDENTITY_ID GRAPHQL_ORCHARD_PROFILE_ID GRAPHQL_ORCHARD_PROFILE_TYPE GRAPHQL_ORCHARD_ROLES Optional env vars: MYSQL_PORT (default: 3306) MYSQL_DATABASE (default: art_relations) """ import argparse import os import sys from dotenv import load_dotenv import pymysql import pymysql.cursors import requests BATCH_SIZE = 500 PRODUCT_LABEL_PARTICIPANTS_QUERY = """ query ProductLabelParticipants($productId: String!) { product(productId: $productId, cached: false) { productId labelParticipations { labelParticipant { ...labelParticipantFragment } } tracks { participations { participant { ...labelParticipantFragment } } labelSoundRecording { participations { role { name } participant { ...labelParticipantFragment } } } } } } fragment labelParticipantFragment on LabelParticipant { uuid name appleMusicId spotifyId spotifyArtistKey globalParticipant { id } } """ CREATE_LABEL_PARTICIPANT_MUTATION = """ mutation CreateLabelParticipant( $input: LabelParticipantInput! $vendorId: Int $subaccountId: Int ) { createLabelParticipant( input: $input vendorId: $vendorId subaccountId: $subaccountId ) { uuid name } } """ def graphql_call(query, variables, vendor_id): headers = { "Content-Type": "application/json", "apollographql-client-name": "jfidlow-move-project", "Orchard-User-Id": os.environ["GRAPHQL_ORCHARD_USER_ID"], "Orchard-Identity-Id": os.environ["GRAPHQL_ORCHARD_IDENTITY_ID"], "Orchard-Profile-Id": os.environ["GRAPHQL_ORCHARD_PROFILE_ID"], "Orchard-Profile-Type": os.environ["GRAPHQL_ORCHARD_PROFILE_TYPE"], "Orchard-Roles": os.environ["GRAPHQL_ORCHARD_ROLES"], "GRASS-ACCOUNT-TYPE": "vendor", "GRASS-ACCOUNT-ID": str(vendor_id), } resp = requests.post( os.environ["GRAPHQL_URL"], json={"query": query, "variables": variables}, headers=headers, timeout=30, ) if not resp.ok: raise RuntimeError( f"GraphQL HTTP {resp.status_code} for vendor {vendor_id}. " f"Response body: {resp.text}" ) body = resp.json() if body.get("errors"): raise RuntimeError(f"GraphQL errors: {body['errors']}") return body["data"] def require_env(*names): missing = [v for v in names if not os.environ.get(v)] if missing: print(f"ERROR: Missing required env vars: {', '.join(missing)}") sys.exit(1) def get_connection(autocommit=True): require_env("MYSQL_HOST", "MYSQL_USER", "MYSQL_PASSWORD") return pymysql.connect( host=os.environ["MYSQL_HOST"], port=int(os.environ.get("MYSQL_PORT", 3306)), user=os.environ["MYSQL_USER"], password=os.environ["MYSQL_PASSWORD"], database=os.environ.get("MYSQL_DATABASE", "art_relations"), autocommit=autocommit, cursorclass=pymysql.cursors.DictCursor, charset="utf8mb4", ) def parse_args(): parser = argparse.ArgumentParser(description="Move a project between vendor accounts") parser.add_argument("--project-id", type=int, required=True) parser.add_argument("--destination-vendor-id", type=int, required=True) parser.add_argument("--destination-subaccount-id", type=int, default=None) parser.add_argument("--dry-run", action="store_true", default=False) return parser.parse_args() def phase1_collect_artists(cursor, project_id): cursor.execute( "SELECT project_id, vendor_id, subaccount_id, artist_id " "FROM project WHERE project_id = %s", (project_id,), ) project_row = cursor.fetchone() if not project_row: print(f"ERROR: project_id {project_id} not found.") sys.exit(1) artist_map = {} proj_artist = project_row["artist_id"] if proj_artist: artist_map[proj_artist] = {"in_project": True, "release_ids": [], "pv_release_ids": []} cursor.execute( "SELECT release_id, artist_id FROM releases WHERE project_id = %s", (project_id,), ) releases = cursor.fetchall() all_release_ids = [] for row in releases: all_release_ids.append(row["release_id"]) aid = row["artist_id"] if aid is None: continue if aid not in artist_map: artist_map[aid] = {"in_project": False, "release_ids": [], "pv_release_ids": []} artist_map[aid]["release_ids"].append(row["release_id"]) if all_release_ids: placeholders = ",".join(["%s"] * len(all_release_ids)) cursor.execute( f"SELECT DISTINCT release_id, primary_artist_id FROM product_video " f"WHERE release_id IN ({placeholders})", tuple(all_release_ids), ) for row in cursor.fetchall(): aid = row["primary_artist_id"] if aid is None: continue if aid not in artist_map: artist_map[aid] = {"in_project": False, "release_ids": [], "pv_release_ids": []} artist_map[aid]["pv_release_ids"].append(row["release_id"]) return project_row, artist_map, all_release_ids def phase2_ensure_artists(cursor, artist_map, dest_vendor_id, dry_run): id_mapping = {} created_count = 0 existed_count = 0 for source_id in artist_map: cursor.execute("SELECT * FROM artist_info WHERE artist_id = %s", (source_id,)) source_row = cursor.fetchone() if not source_row: print(f" WARNING: source artist_id {source_id} not found in artist_info — skipping") continue name = source_row["name"] cursor.execute( "SELECT artist_id FROM artist_info WHERE vendor_id = %s AND name = %s", (dest_vendor_id, name), ) existing = cursor.fetchone() if existing: id_mapping[source_id] = existing["artist_id"] existed_count += 1 continue if dry_run: print(f" [DRY RUN] Would INSERT artist '{name}' for vendor {dest_vendor_id}") id_mapping[source_id] = None created_count += 1 continue cursor.execute( "INSERT IGNORE INTO artist_info (vendor_id, name) VALUES (%s, %s)", (dest_vendor_id, name), ) cursor.execute( "SELECT artist_id FROM artist_info WHERE vendor_id = %s AND name = %s", (dest_vendor_id, name), ) dest_row = cursor.fetchone() id_mapping[source_id] = dest_row["artist_id"] created_count += 1 print( f"Phase 2: {len(artist_map)} unique source artist(s) processed. " f"Created {created_count} new, {existed_count} already existed at destination." ) return id_mapping def phase2_5_capture_and_upsert_label_participants( release_ids, source_vendor_id, dest_vendor_id, dest_subaccount_id, dry_run ): if not release_ids: print("Phase 2.5: no releases — skipping label participant capture.") return {"unique_lps": 0, "upserted": 0} lps_by_uuid = {} for release_id in release_ids: data = graphql_call( PRODUCT_LABEL_PARTICIPANTS_QUERY, {"productId": str(release_id)}, vendor_id=source_vendor_id, ) product = data.get("product") if not product: print(f" WARNING: GraphQL product({release_id}) returned null — skipping.") continue for lp_wrap in product.get("labelParticipations") or []: lp = lp_wrap.get("labelParticipant") if lp and lp.get("uuid"): lps_by_uuid[lp["uuid"]] = lp for track in product.get("tracks") or []: for p_wrap in track.get("participations") or []: lp = p_wrap.get("participant") if lp and lp.get("uuid"): lps_by_uuid[lp["uuid"]] = lp lsr = track.get("labelSoundRecording") if lsr: for p_wrap in lsr.get("participations") or []: lp = p_wrap.get("participant") if lp and lp.get("uuid"): lps_by_uuid[lp["uuid"]] = lp print( f"Phase 2.5: queried {len(release_ids)} release(s); " f"found {len(lps_by_uuid)} unique LabelParticipant(s)." ) subaccount_arg = dest_subaccount_id or 0 if dry_run: print( f" [DRY RUN] Would upsert {len(lps_by_uuid)} unique LabelParticipant(s) " f"on vendor {dest_vendor_id} / subaccount {subaccount_arg}" ) for uuid_, lp in lps_by_uuid.items(): print(f" - {uuid_} {lp.get('name')}") return {"unique_lps": len(lps_by_uuid), "upserted": 0} upserted = 0 for uuid_, lp in lps_by_uuid.items(): input_payload = { "name": lp.get("name"), "appleMusicId": lp.get("appleMusicId"), "spotifyId": lp.get("spotifyId"), "spotifyArtistKey": lp.get("spotifyArtistKey"), } input_payload = {k: v for k, v in input_payload.items() if v is not None} graphql_call( CREATE_LABEL_PARTICIPANT_MUTATION, { "input": input_payload, "vendorId": dest_vendor_id, "subaccountId": subaccount_arg, }, vendor_id=dest_vendor_id, ) upserted += 1 print( f" Upserted {upserted} LabelParticipant(s) on vendor {dest_vendor_id} " f"/ subaccount {subaccount_arg}." ) return {"unique_lps": len(lps_by_uuid), "upserted": upserted} def phase3_execute_transfer(cursor, project_id, project_row, artist_map, id_mapping, dest_vendor_id, dest_subaccount_id, dry_run): stats = {"releases_updated": 0, "pv_rows_updated": 0} proj_artist_src = project_row["artist_id"] proj_artist_dest = id_mapping.get(proj_artist_src) if proj_artist_src else None proj_subaccount = dest_subaccount_id if dest_subaccount_id is not None else 0 if dry_run: print( f" [DRY RUN] Would UPDATE project {project_id}: " f"vendor_id={dest_vendor_id}, subaccount_id={proj_subaccount}, " f"artist_id={proj_artist_dest}" ) else: cursor.execute( "UPDATE project SET vendor_id = %s, subaccount_id = %s, artist_id = %s " "WHERE project_id = %s", (dest_vendor_id, proj_subaccount, proj_artist_dest, project_id), ) for source_id, info in artist_map.items(): dest_id = id_mapping.get(source_id) release_ids = info["release_ids"] pv_release_ids = info["pv_release_ids"] if release_ids: if dry_run: print( f" [DRY RUN] Would UPDATE {len(release_ids)} release(s): " f"artist_id {source_id} → {dest_id}, subaccount_id → {dest_subaccount_id}" ) else: for i in range(0, len(release_ids), BATCH_SIZE): batch = release_ids[i : i + BATCH_SIZE] ph = ",".join(["%s"] * len(batch)) cursor.execute( f"UPDATE releases SET artist_id = %s, subaccount_id = %s " f"WHERE release_id IN ({ph}) AND artist_id = %s", [dest_id, dest_subaccount_id] + batch + [source_id], ) stats["releases_updated"] += cursor.rowcount if pv_release_ids: if dry_run: print( f" [DRY RUN] Would UPDATE product_video for {len(pv_release_ids)} release_id(s): " f"primary_artist_id {source_id} → {dest_id}" ) else: for i in range(0, len(pv_release_ids), BATCH_SIZE): batch = pv_release_ids[i : i + BATCH_SIZE] ph = ",".join(["%s"] * len(batch)) cursor.execute( f"UPDATE product_video SET primary_artist_id = %s " f"WHERE release_id IN ({ph}) AND primary_artist_id = %s", [dest_id] + batch + [source_id], ) stats["pv_rows_updated"] += cursor.rowcount return stats def main(): load_dotenv() args = parse_args() project_id = args.project_id dest_vendor_id = args.destination_vendor_id dest_subaccount_id = args.destination_subaccount_id dry_run = args.dry_run require_env( "MYSQL_HOST", "MYSQL_USER", "MYSQL_PASSWORD", "GRAPHQL_URL", "GRAPHQL_ORCHARD_USER_ID", "GRAPHQL_ORCHARD_IDENTITY_ID", "GRAPHQL_ORCHARD_PROFILE_ID", "GRAPHQL_ORCHARD_PROFILE_TYPE", "GRAPHQL_ORCHARD_ROLES", ) if dry_run: print("=== DRY RUN MODE — no changes will be written ===\n") print("Connecting to database...") conn_ac = get_connection(autocommit=True) try: print(f"\nPhase 1: Collecting artist IDs for project {project_id}...") with conn_ac.cursor() as cur: project_row, artist_map, release_ids = phase1_collect_artists(cur, project_id) if project_row["vendor_id"] == dest_vendor_id: print(f"Project {project_id} is already on vendor {dest_vendor_id}. Nothing to do.") sys.exit(0) print( f" Source vendor: {project_row['vendor_id']}, Destination vendor: {dest_vendor_id}" ) print( f" Found {len(artist_map)} unique artist ID(s) across " f"project + releases + product_video" ) print(f"\nPhase 2: Ensuring all artists exist on destination vendor {dest_vendor_id}...") with conn_ac.cursor() as cur: id_mapping = phase2_ensure_artists(cur, artist_map, dest_vendor_id, dry_run) finally: conn_ac.close() source_vendor_id = project_row["vendor_id"] print( f"\nPhase 2.5: Capturing label participants from {len(release_ids)} release(s) " f"and upserting on destination vendor..." ) phase2_5_capture_and_upsert_label_participants( release_ids, source_vendor_id, dest_vendor_id, dest_subaccount_id, dry_run ) if dry_run: print("\nPhase 3 (DRY RUN): Showing what the transaction would do...") phase3_execute_transfer( None, project_id, project_row, artist_map, id_mapping, dest_vendor_id, dest_subaccount_id, dry_run=True, ) print("\n=== DRY RUN complete. No changes made. ===") return print("\nPhase 3: Executing atomic transfer transaction...") conn_tx = get_connection(autocommit=False) try: with conn_tx.cursor() as cur: stats = phase3_execute_transfer( cur, project_id, project_row, artist_map, id_mapping, dest_vendor_id, dest_subaccount_id, dry_run=False, ) conn_tx.commit() print( f" COMMITTED. Updated: 1 project row, " f"{stats['releases_updated']} release(s), " f"{stats['pv_rows_updated']} product_video row(s)." ) except Exception as e: conn_tx.rollback() print(f"ERROR during transaction — rolled back. Detail: {e}") sys.exit(1) finally: conn_tx.close() print("\nDone.") if __name__ == "__main__": main()