Files
service-finder/backend/app/scripts/migrate_legacy_garages.py

307 lines
12 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 CRITICAL - The Great Garage & Vehicle Migration (Legacy Data Cleanup)
Migrates legacy person_id-based ownership to the new organization_id / organization_members
RBAC structure. Ensures every user has a personal garage (INDIVIDUAL org), is bound as OWNER
in organization_members, and orphaned vehicles are rescued into their owner's garage.
TASKS:
A) Ensure Personal Garages exist for every user who lacks one
B) Bind users to their garages via organization_members (OWNER role) + update scope_id
C) Rescue orphaned vehicles (owner_person_id set but no current_organization_id/owner_org_id)
Usage:
docker exec sf_api python3 /app/app/scripts/migrate_legacy_garages.py
"""
import asyncio
import asyncpg
import json
import os
import sys
from datetime import datetime, timezone
# asyncpg needs plain "postgresql://" scheme, not "postgresql+asyncpg://"
_raw_url = os.getenv(
"DATABASE_URL",
"postgresql://service_finder_app:AppSafePass_2026@db:5432/service_finder"
)
DATABASE_URL = _raw_url.replace("+asyncpg", "")
async def main():
print("=" * 70)
print("🔧 P0 MIGRATION: THE GREAT GARAGE & VEHICLE MIGRATION")
print("=" * 70)
print(f"Started at: {datetime.now(timezone.utc).isoformat()}")
print()
conn = await asyncpg.connect(DATABASE_URL)
try:
# ── TASK A: ENSURE PERSONAL GARAGES EXIST ──
print("📋 TASK A: Ensure Personal Garages Exist")
print("-" * 50)
# Find all users who do NOT have an INDIVIDUAL-type organization where they are owner
users_without_garage = await conn.fetch("""
SELECT u.id AS user_id,
u.email,
p.id AS person_id,
p.first_name,
p.last_name
FROM identity.users u
JOIN identity.persons p ON p.id = u.person_id
WHERE u.is_deleted = false
AND u.is_active = true
AND NOT EXISTS (
SELECT 1
FROM fleet.organizations o
WHERE o.owner_id = u.id
AND o.org_type = 'individual'
AND o.is_deleted = false
)
ORDER BY u.id
""")
print(f" Found {len(users_without_garage)} users without a personal garage.")
garages_created = 0
for row in users_without_garage:
user_id = row['user_id']
person_id = row['person_id']
first_name = row['first_name'] or "Felhasználó"
last_name = row['last_name'] or ""
garage_name = f"{first_name} Garázsa"
folder_slug = f"personal-{user_id}-{person_id}"
now = datetime.now(timezone.utc)
# Direct INSERT with all NOT NULL columns explicitly provided
await conn.execute("""
INSERT INTO fleet.organizations
(full_name, name, display_name, folder_slug, org_type,
owner_id, status, is_active, is_deleted,
subscription_plan, base_asset_limit, purchased_extra_slots,
notification_settings, external_integration_config,
is_ownership_transferable, is_verified,
first_registered_at, current_lifecycle_started_at,
lifecycle_index, created_at,
visual_settings, aliases, tags, settings)
VALUES
($1, $2, $3, $4, 'individual',
$5, 'active', true, false,
'FREE', 5, 0,
$6::json, '{}'::json,
true, false,
$7, $8,
1, $9,
$10::jsonb, '[]'::jsonb, '[]'::jsonb, $11::jsonb)
""",
garage_name, # $1 full_name
garage_name, # $2 name
garage_name, # $3 display_name
folder_slug, # $4 folder_slug
user_id, # $5 owner_id
json.dumps({"notify_owner": True, "alert_days_before": [30, 15, 7, 1]}), # $6 notification_settings
now, # $7 first_registered_at
now, # $8 current_lifecycle_started_at
now, # $9 created_at
json.dumps({"theme": "default", "primary_color": None, "wall_logo_url": None}), # $10 visual_settings
json.dumps({"invite_expiry_days": 7, "max_members": 50, "allow_public_join": False}), # $11 settings
)
garages_created += 1
print(f" ✅ Created personal garage '{garage_name}' for user_id={user_id} ({row['email']})")
print(f"\n 📊 Task A Result: {garages_created} personal garages created.")
print()
# ── TASK B: BIND USERS TO GARAGES (RBAC) ──
print("📋 TASK B: Bind Users to Garages (RBAC)")
print("-" * 50)
# Get all active users with their personal garages
users_and_garages = await conn.fetch("""
SELECT u.id AS user_id,
u.email,
u.scope_id,
o.id AS org_id,
o.name AS org_name
FROM identity.users u
JOIN fleet.organizations o ON o.owner_id = u.id
AND o.org_type = 'individual'
AND o.is_deleted = false
WHERE u.is_deleted = false
AND u.is_active = true
ORDER BY u.id
""")
members_added = 0
scopes_updated = 0
for row in users_and_garages:
user_id = row['user_id']
org_id = row['org_id']
scope_id = row['scope_id']
# Check if user is already a member of this org
existing_member = await conn.fetchrow("""
SELECT id FROM fleet.organization_members
WHERE organization_id = $1 AND user_id = $2
""", org_id, user_id)
if not existing_member:
await conn.execute("""
INSERT INTO fleet.organization_members
(organization_id, user_id, role, status,
is_permanent, is_verified, created_at)
VALUES
($1, $2, 'OWNER', 'active',
true, true, NOW())
""", org_id, user_id)
members_added += 1
print(f" ✅ Added user_id={user_id} as OWNER of org_id={org_id} ({row['org_name']})")
else:
# Ensure role is OWNER
await conn.execute("""
UPDATE fleet.organization_members
SET role = 'OWNER', is_permanent = true, is_verified = true, status = 'active'
WHERE organization_id = $1 AND user_id = $2
""", org_id, user_id)
print(f" Updated membership for user_id={user_id} in org_id={org_id} to OWNER")
# Update scope_id if NULL
if scope_id is None or scope_id == "":
await conn.execute("""
UPDATE identity.users
SET scope_id = $1::text
WHERE id = $2
""", str(org_id), user_id)
scopes_updated += 1
print(f" ✅ Updated scope_id for user_id={user_id} to '{org_id}'")
print(f"\n 📊 Task B Result: {members_added} members added, {scopes_updated} scope_ids updated.")
print()
# ── TASK C: RESCUE ORPHANED VEHICLES ──
print("📋 TASK C: Rescue Orphaned Vehicles")
print("-" * 50)
# Find orphaned vehicles: have owner_person_id but no current_organization_id or owner_org_id
orphaned_vehicles = await conn.fetch("""
SELECT a.id AS asset_id,
a.license_plate,
a.vin,
a.owner_person_id,
a.current_organization_id,
a.owner_org_id,
a.brand,
a.model
FROM vehicle.assets a
WHERE a.owner_person_id IS NOT NULL
AND (a.current_organization_id IS NULL OR a.owner_org_id IS NULL)
AND a.status != 'deleted'
ORDER BY a.owner_person_id, a.id
""")
print(f" Found {len(orphaned_vehicles)} orphaned vehicles.")
vehicles_rescued = 0
vehicles_skipped = 0
for vehicle in orphaned_vehicles:
owner_person_id = vehicle['owner_person_id']
asset_id = vehicle['asset_id']
plate = vehicle['license_plate'] or vehicle['vin'] or str(asset_id)
# Find the user who owns this person record
user_row = await conn.fetchrow("""
SELECT u.id AS user_id
FROM identity.users u
WHERE u.person_id = $1
AND u.is_deleted = false
AND u.is_active = true
LIMIT 1
""", owner_person_id)
if not user_row:
print(f" ⚠️ No active user found for person_id={owner_person_id}, skipping vehicle {plate}")
vehicles_skipped += 1
continue
user_id = user_row['user_id']
# Find the user's personal garage
garage = await conn.fetchrow("""
SELECT o.id AS org_id
FROM fleet.organizations o
WHERE o.owner_id = $1
AND o.org_type = 'individual'
AND o.is_deleted = false
LIMIT 1
""", user_id)
if not garage:
print(f" ⚠️ No personal garage found for user_id={user_id}, skipping vehicle {plate}")
vehicles_skipped += 1
continue
org_id = garage['org_id']
# Update the vehicle's organization fields
updates = []
params = []
param_idx = 1
if vehicle['current_organization_id'] is None:
updates.append(f"current_organization_id = ${param_idx}")
params.append(org_id)
param_idx += 1
if vehicle['owner_org_id'] is None:
updates.append(f"owner_org_id = ${param_idx}")
params.append(org_id)
param_idx += 1
if updates:
params.append(asset_id)
update_sql = f"""
UPDATE vehicle.assets
SET {', '.join(updates)}
WHERE id = ${param_idx}::uuid
"""
await conn.execute(update_sql, *params)
vehicles_rescued += 1
print(f" ✅ Rescued vehicle {plate} → org_id={org_id} (user_id={user_id})")
print(f"\n 📊 Task C Result: {vehicles_rescued} vehicles rescued, {vehicles_skipped} skipped.")
print()
# ── FINAL REPORT ──
print("=" * 70)
print("📊 FINAL MIGRATION REPORT")
print("=" * 70)
print(f" 🏗️ Task A - New Garages Created: {garages_created}")
print(f" 👥 Task B - Members Added: {members_added}")
print(f" 🆔 Task B - Scope IDs Updated: {scopes_updated}")
print(f" 🚗 Task C - Vehicles Rescued: {vehicles_rescued}")
print(f" ⏭️ Task C - Vehicles Skipped: {vehicles_skipped}")
print()
print(f"Completed at: {datetime.now(timezone.utc).isoformat()}")
print("=" * 70)
except Exception as e:
print(f"\n❌ FATAL ERROR: {e}")
import traceback
traceback.print_exc()
sys.exit(1)
finally:
await conn.close()
if __name__ == "__main__":
asyncio.run(main())