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

243 lines
9.6 KiB
Python

"""
🤖 Data Healing Script — Wallet & Referral Code Repair
Detects and fixes missing Wallet and InvitationCode (referral_code) records
for existing User and Organization entities.
Usage:
docker compose exec sf_api python3 /app/backend/app/scripts/heal_user_data.py
Logs:
- Prints a summary of healed records to stdout
- Writes detailed log to logs/heal_user_data_{timestamp}.log
"""
import asyncio
import logging
import os
import sys
import uuid
from datetime import datetime
from typing import Dict, Any
# Ensure the backend app is on the path
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from sqlalchemy import select, text
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine, async_sessionmaker
from app.core.config import settings
from app.models.identity import User, Wallet
from app.models.marketplace.organization import Organization
# ── Logging setup ──────────────────────────────────────────────────────────
LOG_DIR = "/app/backend/logs"
os.makedirs(LOG_DIR, exist_ok=True)
log_filename = f"heal_user_data_{datetime.now().strftime('%Y%m%d_%H%M%S')}.log"
log_path = os.path.join(LOG_DIR, log_filename)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
handlers=[
logging.StreamHandler(sys.stdout),
logging.FileHandler(log_path),
],
)
logger = logging.getLogger("heal-user-data")
# ── Helpers ────────────────────────────────────────────────────────────────
def generate_referral_code() -> str:
"""Generate a short, unique referral code (8 chars, uppercase)."""
import secrets
import string
alphabet = string.ascii_uppercase + string.digits
return "".join(secrets.choice(alphabet) for _ in range(8))
# ── Main healing logic ─────────────────────────────────────────────────────
async def heal_users(db: AsyncSession) -> Dict[str, int]:
"""
Find all Users without a Wallet and create one (0 balance).
Find all Users without a referral_code and generate one.
"""
stats = {"wallet_created": 0, "referral_code_generated": 0}
# 1. Fetch all users
result = await db.execute(select(User))
users = result.scalars().all()
logger.info(f"Found {len(users)} total User records.")
for user in users:
# ── Wallet check ──
wallet_stmt = select(Wallet).where(Wallet.user_id == user.id)
wallet_result = await db.execute(wallet_stmt)
existing_wallet = wallet_result.scalar_one_or_none()
if existing_wallet is None:
new_wallet = Wallet(
user_id=user.id,
earned_credits=0,
purchased_credits=0,
service_coins=0,
currency="HUF",
)
db.add(new_wallet)
stats["wallet_created"] += 1
logger.info(f" [Wallet] Created for User ID={user.id} ({user.email})")
else:
logger.debug(f" [Wallet] Already exists for User ID={user.id}")
# ── Referral code check ──
if not user.referral_code:
code = generate_referral_code()
# Ensure uniqueness
while True:
dup_check = await db.execute(
select(User).where(User.referral_code == code)
)
if dup_check.scalar_one_or_none() is None:
break
code = generate_referral_code()
user.referral_code = code
stats["referral_code_generated"] += 1
logger.info(
f" [Referral] Generated code '{code}' for User ID={user.id} ({user.email})"
)
else:
logger.debug(
f" [Referral] Already has code '{user.referral_code}' for User ID={user.id}"
)
return stats
async def heal_organizations(db: AsyncSession) -> Dict[str, int]:
"""
Check Organizations for wallet coverage.
NOTE: The `wallets` table has:
- user_id: NOT NULL + UNIQUE constraint
- organization_id: nullable
This means each user can have exactly one wallet, and organization
wallets share the user's wallet via organization_id. Since user wallets
are already created in heal_users(), org wallets are inherently covered.
We only log existing org wallet links for audit purposes.
"""
stats = {"org_wallet_created": 0, "org_wallet_skipped": 0}
result = await db.execute(select(Organization))
orgs = result.scalars().all()
logger.info(f"Found {len(orgs)} total Organization records.")
for org in orgs:
# Check if an org wallet already exists (linked via organization_id)
wallet_stmt = select(Wallet).where(Wallet.organization_id == org.id)
wallet_result = await db.execute(wallet_stmt)
existing_wallet = wallet_result.scalar_one_or_none()
if existing_wallet is None:
# The wallets.user_id is NOT NULL + UNIQUE, so we cannot create
# a separate wallet for an org. The org's owner already has a
# personal wallet from heal_users(). We skip org wallet creation
# and log the situation.
owner_id = getattr(org, 'owner_id', None)
if owner_id:
# Check if owner has a wallet we could link
owner_wallet = await db.execute(
select(Wallet).where(Wallet.user_id == owner_id)
)
owner_wallet = owner_wallet.scalar_one_or_none()
if owner_wallet:
# Link the existing wallet to this org
owner_wallet.organization_id = org.id
stats["org_wallet_created"] += 1
logger.info(
f" [OrgWallet] Linked existing wallet (user_id={owner_id}) "
f"to Organization ID={org.id} ({org.name})"
)
else:
stats["org_wallet_skipped"] += 1
logger.warning(
f" [OrgWallet] SKIPPED for Organization ID={org.id} ({org.name}) — "
f"owner (user_id={owner_id}) has no wallet yet."
)
else:
stats["org_wallet_skipped"] += 1
logger.warning(
f" [OrgWallet] SKIPPED for Organization ID={org.id} ({org.name}) — "
f"no owner_id found."
)
else:
logger.debug(f" [OrgWallet] Already linked for Organization ID={org.id}")
return stats
async def main():
"""Main entry point."""
logger.info("=" * 60)
logger.info(" DATA HEALING SCRIPT — Wallet & Referral Code Repair")
logger.info("=" * 60)
# Build async engine from the same DATABASE_URL
db_url = settings.DATABASE_URL
# If the URL starts with postgresql://, convert to postgresql+asyncpg://
if db_url.startswith("postgresql://") and "+asyncpg" not in db_url:
db_url = db_url.replace("postgresql://", "postgresql+asyncpg://", 1)
engine = create_async_engine(db_url, echo=False, pool_pre_ping=True)
async_session_factory = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with async_session_factory() as db:
try:
async with db.begin():
user_stats = await heal_users(db)
org_stats = await heal_organizations(db)
total_healed = (
user_stats["wallet_created"]
+ user_stats["referral_code_generated"]
+ org_stats["org_wallet_created"]
)
logger.info("" * 60)
logger.info(" ✅ HEALING COMPLETE — Summary")
logger.info(f" User wallets created: {user_stats['wallet_created']}")
logger.info(f" Referral codes generated: {user_stats['referral_code_generated']}")
logger.info(f" Organization wallets created: {org_stats['org_wallet_created']}")
logger.info(f" Organization wallets skipped: {org_stats['org_wallet_skipped']}")
logger.info(f" ─────────────────────────────────")
logger.info(f" TOTAL records healed: {total_healed}")
logger.info("" * 60)
print(f"\n{'='*60}")
print(f" ✅ DATA HEALING COMPLETE")
print(f" Log file: {log_path}")
print(f"{'='*60}")
print(f" User wallets created: {user_stats['wallet_created']}")
print(f" Referral codes generated: {user_stats['referral_code_generated']}")
print(f" Organization wallets created: {org_stats['org_wallet_created']}")
print(f" Organization wallets skipped: {org_stats['org_wallet_skipped']}")
print(f" ─────────────────────────────────")
print(f" TOTAL records healed: {total_healed}")
print(f"{'='*60}\n")
except Exception as e:
logger.error(f"❌ Healing failed: {e}", exc_info=True)
print(f"\n❌ ERROR: Healing failed — {e}\n")
raise
await engine.dispose()
if __name__ == "__main__":
asyncio.run(main())