Files
service-finder/backend/scripts/surgical_migrate.py

483 lines
22 KiB
Python
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""
P0 SURGICAL MIGRATION — Target-Only Move to Marketplace
CRITICAL DIRECTIVE: Strict INCLUSION list. ONLY the following Target IDs
will be migrated from fleet.organizations to marketplace.service_providers
and then purged. Every other record in the database remains strictly untouched.
Target IDs: [7676, 4859, 2406, 63, 62, 58]
Folyamat per Target ID (if it exists):
Step A: Read org data → Create ServiceProvider in marketplace.service_providers
Step B: Update fleet_finance.asset_costs WHERE vendor_organization_id == Target ID
→ Set service_provider_id = New_SP_ID, vendor_organization_id = NULL
Step C: Update marketplace.service_profiles WHERE organization_id == Target ID
→ Set service_provider_id = New_SP_ID, organization_id = NULL
Step D (Cleanup): DELETE branches, organization_members, contact_persons, org_relationships
Step E (Final): DELETE the Target ID from fleet.organizations
Használat:
docker compose exec sf_api python3 /app/backend/scripts/surgical_migrate.py
"""
import asyncio
import logging
import sys
from datetime import datetime, timezone
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
sys.path.insert(0, '/app')
from app.core.config import settings
# ──────────────────────────────────────────────
# Konfiguráció
# ──────────────────────────────────────────────
# STRICT INCLUSION LIST — ONLY these IDs will be touched
TARGET_ORG_IDS = [7676, 4859, 2406, 63, 62, 58]
DATABASE_URL = settings.DATABASE_URL
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
handlers=[logging.StreamHandler(sys.stdout)],
)
logger = logging.getLogger(__name__)
async def surgical_migrate():
"""
Surgical migration — ONLY iterates through the explicitly provided Target IDs.
"""
engine = create_async_engine(DATABASE_URL, echo=False)
async with AsyncSession(engine) as db:
try:
# =====================================================================
# PHASE 0: Verify which target IDs actually exist
# =====================================================================
logger.info("=" * 60)
logger.info("PHASE 0: Target ID Verification")
logger.info("=" * 60)
target_ids_str = ", ".join([str(x) for x in TARGET_ORG_IDS])
verify_sql = f"""
SELECT id, name, org_type, status, is_deleted
FROM fleet.organizations
WHERE id IN ({target_ids_str})
ORDER BY id
"""
result = await db.execute(text(verify_sql))
existing_orgs = result.fetchall()
existing_ids = {row[0] for row in existing_orgs}
missing_ids = [oid for oid in TARGET_ORG_IDS if oid not in existing_ids]
logger.info(f"Target IDs to process: {TARGET_ORG_IDS}")
logger.info(f"Found in database: {len(existing_orgs)} organizations")
for row in existing_orgs:
logger.info(f" ✅ ID={row[0]}: name='{row[1]}', type='{row[2]}', status='{row[3]}', is_deleted={row[4]}")
if missing_ids:
logger.warning(f" ⚠️ Missing IDs (not found in fleet.organizations): {missing_ids}")
if not existing_orgs:
logger.info(" No target organizations found. Nothing to do.")
return {"status": "no_targets_found", "missing_ids": missing_ids}
# =====================================================================
# PHASE 1: STEP A — Migrate each target org to ServiceProvider
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("PHASE 1: STEP A — Migrate to ServiceProvider")
logger.info("=" * 60)
# Track the mapping: org_id -> new_sp_id
org_to_sp_map = {}
for org_row in existing_orgs:
org_id = org_row[0]
org_name = org_row[1]
org_type = org_row[2]
logger.info(f"\n📦 Processing org ID={org_id} ('{org_name}')...")
# Build address from org fields
address_parts = []
if org_row[4]: # We need to re-fetch with more columns
pass
# Re-fetch full org data
full_org_sql = text("""
SELECT id, name, full_name, address_zip, address_city,
address_street_name, address_street_type, address_house_number,
plus_code, tax_number, reg_number, org_type, status
FROM fleet.organizations
WHERE id = :org_id
""")
full_result = await db.execute(full_org_sql, {"org_id": org_id})
org_data = full_result.fetchone()
if not org_data:
logger.warning(f" ⚠️ Org ID={org_id} disappeared between checks. Skipping.")
continue
# Map column indices
col_name = 1
col_full_name = 2
col_zip = 3
col_city = 4
col_street_name = 5
col_street_type = 6
col_house_number = 7
col_plus_code = 8
col_tax_number = 9
col_reg_number = 10
col_org_type = 11
col_status = 12
# Build a meaningful address string
addr_parts = []
if org_data[col_zip]:
addr_parts.append(org_data[col_zip])
if org_data[col_city]:
addr_parts.append(org_data[col_city])
if org_data[col_street_name]:
street = org_data[col_street_name]
if org_data[col_street_type]:
street = f"{org_data[col_street_type]} {street}"
if org_data[col_house_number]:
street = f"{street} {org_data[col_house_number]}"
addr_parts.append(street)
address = ", ".join(addr_parts) if addr_parts else org_data[col_name]
# Infer category from org_type
category_map = {
"service": "service",
"service_provider": "service",
"fleet_owner": "fleet",
"business": "business",
"club": "club",
"individual": "individual",
}
category = category_map.get(org_data[col_org_type], "service")
# Check if a ServiceProvider with this name already exists
check_sp_sql = text("""
SELECT id FROM marketplace.service_providers
WHERE name = :org_name
LIMIT 1
""")
sp_check = await db.execute(check_sp_sql, {"org_name": org_data[col_name]})
existing_sp = sp_check.fetchone()
if existing_sp:
sp_id = existing_sp[0]
logger.info(f" ServiceProvider already exists for '{org_data[col_name]}' (ID={sp_id}). Reusing.")
else:
# Create new ServiceProvider
insert_sp_sql = text("""
INSERT INTO marketplace.service_providers
(name, address, category, city, address_zip,
address_street_name, address_street_type, address_house_number,
plus_code, status, source, validation_score, created_at)
VALUES
(:name, :address, :category, :city, :zip,
:street_name, :street_type, :house_number,
:plus_code, 'approved', 'api', 100, NOW())
RETURNING id
""")
sp_params = {
"name": org_data[col_name],
"address": address,
"category": category,
"city": org_data[col_city],
"zip": org_data[col_zip],
"street_name": org_data[col_street_name],
"street_type": org_data[col_street_type],
"house_number": org_data[col_house_number],
"plus_code": org_data[col_plus_code],
}
sp_result = await db.execute(insert_sp_sql, sp_params)
sp_id = sp_result.scalar()
logger.info(f" ✅ Created ServiceProvider ID={sp_id} for org '{org_data[1]}'")
org_to_sp_map[org_id] = sp_id
if not org_to_sp_map:
logger.warning("⚠️ No organizations were migrated. Aborting.")
await db.rollback()
return {"status": "no_migrations_performed"}
# =====================================================================
# PHASE 2: STEP B — Relink AssetCost records
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("PHASE 2: STEP B — Relink AssetCost (vendor_organization_id -> service_provider_id)")
logger.info("=" * 60)
total_ac_relinked = 0
total_ac_skipped = 0
total_ac_errors = 0
for org_id, sp_id in org_to_sp_map.items():
logger.info(f"\n📊 Processing AssetCosts for org ID={org_id} -> SP ID={sp_id}...")
# Find AssetCost records with this vendor_organization_id
ac_find_sql = text("""
SELECT id, vendor_organization_id, service_provider_id
FROM fleet_finance.asset_costs
WHERE vendor_organization_id = :org_id
""")
ac_result = await db.execute(ac_find_sql, {"org_id": org_id})
asset_costs = ac_result.fetchall()
if not asset_costs:
logger.info(f" No AssetCost records found for vendor_organization_id={org_id}")
continue
for ac_row in asset_costs:
ac_id = ac_row[0]
try:
update_ac_sql = text("""
UPDATE fleet_finance.asset_costs
SET service_provider_id = :sp_id,
vendor_organization_id = NULL
WHERE id = :ac_id
""")
await db.execute(update_ac_sql, {"sp_id": sp_id, "ac_id": ac_id})
total_ac_relinked += 1
logger.info(f" ✅ AssetCost {ac_id}: vendor_org={org_id} -> service_provider_id={sp_id}")
except Exception as e:
total_ac_errors += 1
logger.error(f" ❌ Error relinking AssetCost {ac_id}: {e}")
logger.info(f"\n📊 AssetCost Summary: {total_ac_relinked} relinked, {total_ac_errors} errors")
# =====================================================================
# PHASE 3: STEP C — Relink ServiceProfile records
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("PHASE 3: STEP C — Relink ServiceProfile (organization_id -> service_provider_id)")
logger.info("=" * 60)
total_profiles_relinked = 0
total_profiles_skipped = 0
total_profiles_errors = 0
for org_id, sp_id in org_to_sp_map.items():
logger.info(f"\n📊 Processing ServiceProfiles for org ID={org_id} -> SP ID={sp_id}...")
prof_find_sql = text("""
SELECT id, organization_id, service_provider_id
FROM marketplace.service_profiles
WHERE organization_id = :org_id
""")
prof_result = await db.execute(prof_find_sql, {"org_id": org_id})
profiles = prof_result.fetchall()
if not profiles:
logger.info(f" No ServiceProfile records found for organization_id={org_id}")
continue
for prof_row in profiles:
profile_id = prof_row[0]
try:
update_prof_sql = text("""
UPDATE marketplace.service_profiles
SET service_provider_id = :sp_id,
organization_id = NULL
WHERE id = :profile_id
""")
await db.execute(update_prof_sql, {"sp_id": sp_id, "profile_id": profile_id})
total_profiles_relinked += 1
logger.info(f" ✅ ServiceProfile {profile_id}: org={org_id} -> service_provider_id={sp_id}")
except Exception as e:
total_profiles_errors += 1
logger.error(f" ❌ Error relinking ServiceProfile {profile_id}: {e}")
logger.info(f"\n📊 ServiceProfile Summary: {total_profiles_relinked} relinked, {total_profiles_errors} errors")
# =====================================================================
# PHASE 4: STEP D — Cleanup (Delete dependent records)
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("PHASE 4: STEP D — Cleanup Dependent Records")
logger.info("=" * 60)
org_ids_for_cleanup = list(org_to_sp_map.keys())
org_ids_cleanup_str = ", ".join([str(x) for x in org_ids_for_cleanup])
# D1. Delete branches
logger.info("\n🗑️ D1. Deleting branches...")
del_branches_sql = f"""
DELETE FROM fleet.branches
WHERE organization_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_branches_sql))
deleted_branches = result.rowcount
logger.info(f" ✅ Deleted {deleted_branches} branches")
# D2. Delete organization_members
logger.info("\n🗑️ D2. Deleting organization_members...")
del_members_sql = f"""
DELETE FROM fleet.organization_members
WHERE organization_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_members_sql))
deleted_members = result.rowcount
logger.info(f" ✅ Deleted {deleted_members} organization members")
# D3. Delete contact_persons
logger.info("\n🗑️ D3. Deleting contact_persons...")
del_contacts_sql = f"""
DELETE FROM fleet.contact_persons
WHERE organization_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_contacts_sql))
deleted_contacts = result.rowcount
logger.info(f" ✅ Deleted {deleted_contacts} contact persons")
# D4. Delete org_relationships (both source and target)
logger.info("\n🗑️ D4. Deleting org_relationships...")
del_rels_sql = f"""
DELETE FROM fleet.org_relationships
WHERE source_org_id IN ({org_ids_cleanup_str})
OR target_org_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_rels_sql))
deleted_rels = result.rowcount
logger.info(f" ✅ Deleted {deleted_rels} org relationships")
# D5. Delete organization_financials
logger.info("\n🗑️ D5. Deleting organization_financials...")
del_financials_sql = f"""
DELETE FROM fleet.organization_financials
WHERE organization_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_financials_sql))
deleted_financials = result.rowcount
logger.info(f" ✅ Deleted {deleted_financials} organization financials")
# D6. Delete org_sales_assignments
logger.info("\n🗑️ D6. Deleting org_sales_assignments...")
try:
del_sales_sql = f"""
DELETE FROM fleet.org_sales_assignments
WHERE organization_id IN ({org_ids_cleanup_str})
"""
result = await db.execute(text(del_sales_sql))
deleted_sales = result.rowcount
logger.info(f" ✅ Deleted {deleted_sales} sales assignments")
except Exception as e:
logger.warning(f" ⚠️ Error deleting sales assignments: {e}")
deleted_sales = 0
# =====================================================================
# PHASE 5: STEP E — Final Purge: Delete the organizations
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("PHASE 5: STEP E — Final Purge: Delete organizations")
logger.info("=" * 60)
del_orgs_sql = f"""
DELETE FROM fleet.organizations
WHERE id IN ({org_ids_cleanup_str})
RETURNING id, name
"""
result = await db.execute(text(del_orgs_sql))
deleted_orgs = result.fetchall()
deleted_orgs_count = len(deleted_orgs)
logger.info(f"\n💀 Deleted {deleted_orgs_count} organizations:")
for org_id, org_name in deleted_orgs:
logger.info(f" - ID={org_id}: {org_name}")
# ── Commit ──
await db.commit()
# =====================================================================
# FINAL REPORT
# =====================================================================
logger.info("\n" + "=" * 60)
logger.info("📋 FINAL SURGICAL MIGRATION REPORT")
logger.info("=" * 60)
logger.info(f"\n🎯 Target IDs: {TARGET_ORG_IDS}")
if missing_ids:
logger.info(f"⚠️ Missing IDs (not found): {missing_ids}")
logger.info(f"✅ Successfully migrated: {list(org_to_sp_map.keys())}")
logger.info(f"\n📊 ServiceProvider Mapping (org_id -> sp_id):")
for org_id, sp_id in org_to_sp_map.items():
logger.info(f" Org {org_id} -> ServiceProvider {sp_id}")
logger.info(f"\n📊 Relink Summary:")
logger.info(f" AssetCost records relinked: {total_ac_relinked}")
logger.info(f" AssetCost errors: {total_ac_errors}")
logger.info(f" ServiceProfile records relinked: {total_profiles_relinked}")
logger.info(f" ServiceProfile errors: {total_profiles_errors}")
logger.info(f"\n🗑️ Cleanup Summary:")
logger.info(f" Branches deleted: {deleted_branches}")
logger.info(f" Organization Members deleted: {deleted_members}")
logger.info(f" Contact Persons deleted: {deleted_contacts}")
logger.info(f" Org Relationships deleted: {deleted_rels}")
logger.info(f" Organization Financials deleted: {deleted_financials}")
logger.info(f" Sales Assignments deleted: {deleted_sales}")
logger.info(f"\n💀 Organizations Purged: {deleted_orgs_count}")
total_garbage = (
deleted_branches + deleted_members + deleted_contacts
+ deleted_rels + deleted_financials + deleted_sales
)
logger.info(f"\n📈 Total garbage rows destroyed: {total_garbage}")
logger.info(f"\n✅ Surgical migration completed successfully!")
return {
"status": "success",
"target_ids": TARGET_ORG_IDS,
"missing_ids": missing_ids,
"migrated_ids": list(org_to_sp_map.keys()),
"org_to_sp_map": {str(k): v for k, v in org_to_sp_map.items()},
"asset_cost_relinked": total_ac_relinked,
"asset_cost_errors": total_ac_errors,
"profiles_relinked": total_profiles_relinked,
"profiles_errors": total_profiles_errors,
"deleted_branches": deleted_branches,
"deleted_members": deleted_members,
"deleted_contacts": deleted_contacts,
"deleted_rels": deleted_rels,
"deleted_financials": deleted_financials,
"deleted_sales": deleted_sales,
"deleted_orgs": deleted_orgs_count,
"total_garbage": total_garbage,
}
except Exception as e:
await db.rollback()
logger.error(f"💥 Fatal error during surgical migration: {e}", exc_info=True)
sys.exit(1)
async def main():
logger.info("🚀 P0 SURGICAL MIGRATION — Target-Only Move to Marketplace")
logger.info(f" Target IDs (INCLUSION LIST): {TARGET_ORG_IDS}")
logger.info(f" Database: {DATABASE_URL}")
logger.info("")
try:
await surgical_migrate()
except Exception as e:
logger.error(f"💥 Fatal error: {e}")
sys.exit(1)
if __name__ == "__main__":
asyncio.run(main())