364 lines
12 KiB
Python
364 lines
12 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
📦 Nagy Előfizetés Migráció (Subscription Migration Script)
|
|
|
|
Cél:
|
|
A régi, denormalizált subscription_plan mezők (identity.users, fleet.organizations)
|
|
átvezetése a modern, normalizált subscription_tier + org_subscriptions / user_subscriptions
|
|
architektúrába.
|
|
|
|
Feladatok:
|
|
A) Létrehozza a corp_free_v1 csomagot, ha még nem létezik (0 áras, alap korlátokkal).
|
|
B) Létrehozza a user_subscriptions táblát (finance.user_subscriptions), ha még nem létezik.
|
|
C) Adatkonverzió: a meglévő subscription_plan string-eket átalakítja tier_id hivatkozásokká.
|
|
|
|
Használat:
|
|
docker compose exec sf_api python3 /app/backend/app/scripts/migrate_subscriptions.py
|
|
|
|
Architektúra:
|
|
- system.subscription_tiers: A csomagdefiníciók (SubscriptionTier modell)
|
|
- finance.org_subscriptions: Szervezeti előfizetések (OrganizationSubscription modell)
|
|
- finance.user_subscriptions: Felhasználói előfizetések (UserSubscription modell - ÚJ!)
|
|
"""
|
|
|
|
import asyncio
|
|
import logging
|
|
from datetime import datetime, timezone
|
|
from sqlalchemy import text, select, Integer, String, Boolean, DateTime, ForeignKey
|
|
from sqlalchemy.orm import Mapped, mapped_column
|
|
from sqlalchemy.sql import func
|
|
|
|
from app.database import AsyncSessionLocal, Base
|
|
from app.models.core_logic import SubscriptionTier, OrganizationSubscription
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(levelname)s] %(message)s",
|
|
)
|
|
logger = logging.getLogger("Migrate-Subscriptions")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tiers név -> ID mapping (feltöltés után)
|
|
# ---------------------------------------------------------------------------
|
|
TIER_MAP: dict[str, int] = {}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# A) corp_free_v1 definíció
|
|
# ---------------------------------------------------------------------------
|
|
CORP_FREE_V1 = {
|
|
"name": "corp_free_v1",
|
|
"rules": {
|
|
"type": "corporate",
|
|
"display_name": "Céges Ingyenes",
|
|
"pricing": {
|
|
"monthly_price": 0.0,
|
|
"yearly_price": 0.0,
|
|
"currency": "EUR",
|
|
"credit_price": 0,
|
|
},
|
|
"allowances": {
|
|
"max_vehicles": 3,
|
|
"max_garages": 1,
|
|
"monthly_free_credits": 0,
|
|
},
|
|
"entitlements": [],
|
|
"affiliate": {
|
|
"commission_rate_percent": 0,
|
|
"referral_bonus_credits": 0,
|
|
},
|
|
"lifecycle": {
|
|
"is_public": True,
|
|
},
|
|
},
|
|
"is_custom": False,
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Régi plan string -> tier név mapping
|
|
# ---------------------------------------------------------------------------
|
|
PLAN_TO_TIER: dict[str, str] = {
|
|
"FREE": "private_free_v1",
|
|
"free": "private_free_v1",
|
|
"PRO": "private_pro_v1",
|
|
"PREMIUM": "corp_premium_v1",
|
|
"ENTERPRISE": "corp_vip_v1",
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Segédfüggvények
|
|
# ---------------------------------------------------------------------------
|
|
async def _ensure_corp_free_v1(db) -> None:
|
|
"""A) Létrehozza a corp_free_v1 csomagot, ha még nem létezik."""
|
|
result = await db.execute(
|
|
select(SubscriptionTier).where(SubscriptionTier.name == "corp_free_v1")
|
|
)
|
|
existing = result.scalar_one_or_none()
|
|
|
|
if existing:
|
|
logger.info("✅ corp_free_v1 már létezik (ID=%d). Kihagyva.", existing.id)
|
|
TIER_MAP["corp_free_v1"] = existing.id
|
|
return
|
|
|
|
tier = SubscriptionTier(
|
|
name=CORP_FREE_V1["name"],
|
|
rules=CORP_FREE_V1["rules"],
|
|
is_custom=CORP_FREE_V1["is_custom"],
|
|
)
|
|
db.add(tier)
|
|
await db.commit()
|
|
await db.refresh(tier)
|
|
TIER_MAP["corp_free_v1"] = tier.id
|
|
logger.info("✅ corp_free_v1 létrehozva (ID=%d).", tier.id)
|
|
|
|
|
|
async def _load_tier_map(db) -> None:
|
|
"""Betölti az összes tier nevét és ID-ját a TIER_MAP-ba."""
|
|
result = await db.execute(select(SubscriptionTier))
|
|
tiers = result.scalars().all()
|
|
for t in tiers:
|
|
TIER_MAP[t.name] = t.id
|
|
logger.info("📋 Betöltött tier-ek: %s", {n: i for n, i in TIER_MAP.items()})
|
|
|
|
|
|
def _resolve_tier_id(plan: str | None) -> int | None:
|
|
"""Egy régi plan string-ből kikeresi a megfelelő tier ID-t."""
|
|
if not plan:
|
|
return None
|
|
tier_name = PLAN_TO_TIER.get(plan)
|
|
if not tier_name:
|
|
logger.warning("⚠️ Ismeretlen plan: '%s' -> private_free_v1 lesz.", plan)
|
|
tier_name = "private_free_v1"
|
|
return TIER_MAP.get(tier_name)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# B) UserSubscription tábla létrehozása (ha nem létezik)
|
|
# ---------------------------------------------------------------------------
|
|
async def _ensure_user_subscriptions_table(db) -> None:
|
|
"""Létrehozza a finance.user_subscriptions táblát, ha még nem létezik."""
|
|
result = await db.execute(
|
|
text(
|
|
"SELECT EXISTS (SELECT FROM information_schema.tables "
|
|
"WHERE table_schema = 'finance' AND table_name = 'user_subscriptions')"
|
|
)
|
|
)
|
|
exists = result.scalar()
|
|
|
|
if exists:
|
|
logger.info("✅ finance.user_subscriptions tábla már létezik.")
|
|
return
|
|
|
|
logger.info("🛠️ Létrehozom a finance.user_subscriptions táblát...")
|
|
await db.execute(
|
|
text("""
|
|
CREATE TABLE finance.user_subscriptions (
|
|
id SERIAL PRIMARY KEY,
|
|
user_id INTEGER NOT NULL REFERENCES identity.users(id) ON DELETE CASCADE,
|
|
tier_id INTEGER NOT NULL REFERENCES system.subscription_tiers(id),
|
|
valid_from TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
valid_until TIMESTAMPTZ,
|
|
is_active BOOLEAN NOT NULL DEFAULT TRUE,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ
|
|
)
|
|
""")
|
|
)
|
|
# Index a gyors lekérdezéshez
|
|
await db.execute(
|
|
text("""
|
|
CREATE INDEX idx_user_subscriptions_user_id
|
|
ON finance.user_subscriptions(user_id)
|
|
""")
|
|
)
|
|
await db.commit()
|
|
logger.info("✅ finance.user_subscriptions tábla létrehozva.")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# C) Adatkonverzió
|
|
# ---------------------------------------------------------------------------
|
|
async def _migrate_organizations(db) -> None:
|
|
"""Szervezetek subscription_plan mezőjének átvezetése org_subscriptions-ba."""
|
|
logger.info("=" * 60)
|
|
logger.info("Szervezetek migrálása...")
|
|
logger.info("=" * 60)
|
|
|
|
# 1. Lekérdezzük az összes szervezetet, aminek nincs még org_subscription-e
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT o.id, o.subscription_plan
|
|
FROM fleet.organizations o
|
|
WHERE o.is_deleted = FALSE
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM finance.org_subscriptions os
|
|
WHERE os.org_id = o.id
|
|
)
|
|
ORDER BY o.id
|
|
""")
|
|
)
|
|
orgs = result.all()
|
|
logger.info("📊 %d szervezet vár migrálásra.", len(orgs))
|
|
|
|
inserted = 0
|
|
skipped = 0
|
|
for org in orgs:
|
|
tier_id = _resolve_tier_id(org.subscription_plan)
|
|
if tier_id is None:
|
|
logger.warning("⚠️ Org ID=%d: plan='%s' nem feloldható, kihagyva.", org.id, org.subscription_plan)
|
|
skipped += 1
|
|
continue
|
|
|
|
sub = OrganizationSubscription(
|
|
org_id=org.id,
|
|
tier_id=tier_id,
|
|
valid_from=datetime.now(timezone.utc),
|
|
valid_until=None,
|
|
is_active=True,
|
|
)
|
|
db.add(sub)
|
|
inserted += 1
|
|
|
|
await db.commit()
|
|
logger.info("✅ %d org_subscription beszúrva. (%d kihagyva)", inserted, skipped)
|
|
|
|
|
|
async def _migrate_users(db) -> None:
|
|
"""Felhasználók subscription_plan mezőjének átvezetése user_subscriptions-ba."""
|
|
logger.info("=" * 60)
|
|
logger.info("Felhasználók migrálása...")
|
|
logger.info("=" * 60)
|
|
|
|
# 1. Lekérdezzük az összes aktív felhasználót, aminek nincs még user_subscription-e
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT u.id, u.subscription_plan
|
|
FROM identity.users u
|
|
WHERE u.is_deleted = FALSE
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM finance.user_subscriptions us
|
|
WHERE us.user_id = u.id
|
|
)
|
|
ORDER BY u.id
|
|
""")
|
|
)
|
|
users = result.all()
|
|
logger.info("📊 %d felhasználó vár migrálásra.", len(users))
|
|
|
|
inserted = 0
|
|
skipped = 0
|
|
for user in users:
|
|
tier_id = _resolve_tier_id(user.subscription_plan)
|
|
if tier_id is None:
|
|
logger.warning("⚠️ User ID=%d: plan='%s' nem feloldható, kihagyva.", user.id, user.subscription_plan)
|
|
skipped += 1
|
|
continue
|
|
|
|
await db.execute(
|
|
text("""
|
|
INSERT INTO finance.user_subscriptions
|
|
(user_id, tier_id, valid_from, is_active, created_at)
|
|
VALUES
|
|
(:uid, :tid, NOW(), TRUE, NOW())
|
|
"""),
|
|
{"uid": user.id, "tid": tier_id},
|
|
)
|
|
inserted += 1
|
|
|
|
await db.commit()
|
|
logger.info("✅ %d user_subscription beszúrva. (%d kihagyva)", inserted, skipped)
|
|
|
|
|
|
async def _verify_migration(db) -> None:
|
|
"""Ellenőrzi a migráció sikerességét."""
|
|
logger.info("=" * 60)
|
|
logger.info("🔍 Ellenőrzés")
|
|
logger.info("=" * 60)
|
|
|
|
# Tiers
|
|
result = await db.execute(
|
|
text("SELECT id, name FROM system.subscription_tiers ORDER BY id")
|
|
)
|
|
tiers = result.all()
|
|
logger.info("📋 Tiers (%d db):", len(tiers))
|
|
for t in tiers:
|
|
logger.info(" ID=%d name=%s", t.id, t.name)
|
|
|
|
# Org subscriptions
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT COUNT(*) as cnt,
|
|
COALESCE(t.name, 'N/A') as tier_name
|
|
FROM finance.org_subscriptions os
|
|
LEFT JOIN system.subscription_tiers t ON t.id = os.tier_id
|
|
GROUP BY t.name
|
|
ORDER BY t.name
|
|
""")
|
|
)
|
|
rows = result.all()
|
|
total_org = sum(r.cnt for r in rows)
|
|
logger.info("📊 Org subscriptions (%d db):", total_org)
|
|
for r in rows:
|
|
logger.info(" %s: %d", r.tier_name, r.cnt)
|
|
|
|
# User subscriptions
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT COUNT(*) as cnt,
|
|
COALESCE(t.name, 'N/A') as tier_name
|
|
FROM finance.user_subscriptions us
|
|
LEFT JOIN system.subscription_tiers t ON t.id = us.tier_id
|
|
GROUP BY t.name
|
|
ORDER BY t.name
|
|
""")
|
|
)
|
|
rows = result.all()
|
|
total_user = sum(r.cnt for r in rows)
|
|
logger.info("📊 User subscriptions (%d db):", total_user)
|
|
for r in rows:
|
|
logger.info(" %s: %d", r.tier_name, r.cnt)
|
|
|
|
logger.info("=" * 60)
|
|
logger.info("✅ Migráció befejeződött!")
|
|
logger.info(" Összes org subscription: %d", total_org)
|
|
logger.info(" Összes user subscription: %d", total_user)
|
|
logger.info("=" * 60)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Main
|
|
# ---------------------------------------------------------------------------
|
|
async def migrate():
|
|
logger.info("=" * 60)
|
|
logger.info("📦 Nagy Előfizetés Migráció")
|
|
logger.info("=" * 60)
|
|
|
|
async with AsyncSessionLocal() as db:
|
|
# A) corp_free_v1 létrehozása
|
|
logger.info("\n[1/5] corp_free_v1 csomag ellenőrzése...")
|
|
await _ensure_corp_free_v1(db)
|
|
|
|
# Tier map betöltése
|
|
logger.info("\n[2/5] Tier-ek betöltése...")
|
|
await _load_tier_map(db)
|
|
|
|
# B) user_subscriptions tábla létrehozása
|
|
logger.info("\n[3/5] user_subscriptions tábla ellenőrzése...")
|
|
await _ensure_user_subscriptions_table(db)
|
|
|
|
# C) Adatkonverzió
|
|
logger.info("\n[4/5] Adatkonverzió...")
|
|
await _migrate_organizations(db)
|
|
await _migrate_users(db)
|
|
|
|
# Ellenőrzés
|
|
logger.info("\n[5/5] Ellenőrzés...")
|
|
await _verify_migration(db)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(migrate())
|