Several services (ipg, settlement, advertising, promotions) used to pass a user wallet (rial/reward) as their company-side wallet because they had no dedicated company account. They now each do. This adds a one-shot management command that repoints the company side of historical Transaction rows onto the new dedicated wallets and corrects the two affected Account.balance running totals. - apps/wallet/management/commands/backfill_company_wallets.py dry-run by default; --execute; --service <name>; --app-<svc> <uuid> override. Idempotent (filters on the old account), single atomic + select_for_update. - docs/company_wallet_history_backfill.md — full write-up: model, mapping, balance-correction logic, consequences, run checklist, source commits. Not yet run against staging/prod. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
234 lines
11 KiB
Python
234 lines
11 KiB
Python
"""Backfill historical Transaction rows onto the new dedicated company wallets.
|
|
|
|
Several services used to pass a *user* wallet (rial / reward) as the company-side
|
|
account of their wallet calls, because they had no dedicated company account.
|
|
They now each have one. This command repoints the company side of every historical
|
|
``Transaction`` from the old (user) wallet to the new dedicated wallet, and fixes
|
|
the two denormalised ``Account.balance`` running totals so nothing is lost.
|
|
|
|
python manage.py backfill_company_wallets # dry-run, all services
|
|
python manage.py backfill_company_wallets --service ipg # dry-run, one service
|
|
python manage.py backfill_company_wallets --execute # actually write
|
|
|
|
Run it once, in the environment that owns the wallet DB. It is idempotent — a
|
|
second run finds nothing left to move (the filter is on the *old* account).
|
|
|
|
App identities are resolved from the DB at runtime (by Application.name, with the
|
|
new dedicated wallet's existing account as a cross-check). Override with
|
|
--app-<service> <uuid> if resolution is ambiguous; --dry-run prints what it found.
|
|
|
|
Full write-up: docs/company_wallet_history_backfill.md
|
|
"""
|
|
|
|
from django.core.management.base import BaseCommand, CommandError
|
|
from django.db import transaction
|
|
from django.db.models import Q, Sum
|
|
|
|
from apps.gooyal_oauth2.models import Application
|
|
from apps.wallet.constans import StateChoices, TypeChoices
|
|
from apps.wallet.models import Account, Transaction, Wallet
|
|
|
|
# ─────────────────────────────────────────────────────────────────────────────
|
|
# Wallet UUIDs (unchanged user-side wallets + the new dedicated company wallets)
|
|
RIAL = "af7d967f-30c0-409b-9066-2549f2da5e5e" # WALLET_RIAL == WALLET_RIAL_DEPOSIT
|
|
REWARD = "e7c9d4d1-4d1f-43b2-96f7-4d4a168f480d" # WALLET_REWARD == WALLET_USER_BILLBOARD_VISIT_INCOME
|
|
|
|
IPG_CREDIT = "5c693c93-6b13-476e-a720-e38f3798acae"
|
|
SETTLEMENT_TRANSIT = "939d9d70-3bda-4413-9f9e-756ef4e1525a"
|
|
SETTLEMENT_COMMISSION_INCOME = "ee8b050a-0ab7-48c3-a13c-01733de9bb1d"
|
|
ADVERTISING_TRANSIT = "052d38f0-d4de-40ff-85f6-9ee6e880b7e4"
|
|
PROMOTIONS_CREDIT = "f1f14c34-7e28-4d28-97c6-2bb8b2189ff3" # env still named WALLET_PROMOTIONS_TRANSIT
|
|
|
|
# States in which the company-side balance effect is currently applied and must
|
|
# therefore travel with the rows when we repoint them:
|
|
# deposit flow → company is the PAYER, debited at submit() (PENDING) through SUCCESS
|
|
# withdraw flow → company is the PAYEE, credited only at verify() (SUCCESS)
|
|
PAYER_LIVE_STATES = (StateChoices.PENDING, StateChoices.SUCCESS, StateChoices.DELAYED, StateChoices.INCOMPLETE)
|
|
PAYEE_LIVE_STATES = (StateChoices.SUCCESS,)
|
|
|
|
# Persian msgstr variants of the settlement leg descriptions (LANGUAGE_CODE is
|
|
# en-us, but a request-context call may have emitted the translated string).
|
|
COMMISSION_DESCRIPTIONS = ["settlement commission transaction", "کارمزد تسویه حساب"]
|
|
|
|
# ─────────────────────────────────────────────────────────────────────────────
|
|
# service → what to move. Each "leg":
|
|
# side: 'payer' (deposit calls) or 'payee' (withdraw calls) — the company side
|
|
# old / new: wallet UUIDs
|
|
# desc_any: only rows whose details.description matches one of these (optional)
|
|
# desc_none: skip rows whose details.description matches one of these (optional)
|
|
SPECS = {
|
|
"ipg": {
|
|
"app_names": ["ipg", "ipg app", "gateway"],
|
|
"cross_check_wallet": IPG_CREDIT,
|
|
"legs": [
|
|
{"side": "payer", "old": RIAL, "new": IPG_CREDIT},
|
|
],
|
|
},
|
|
"settlement": {
|
|
"app_names": ["settlement"],
|
|
"cross_check_wallet": SETTLEMENT_TRANSIT,
|
|
"legs": [
|
|
{"side": "payee", "old": RIAL, "new": SETTLEMENT_COMMISSION_INCOME, "desc_any": COMMISSION_DESCRIPTIONS},
|
|
{"side": "payee", "old": REWARD, "new": SETTLEMENT_COMMISSION_INCOME, "desc_any": COMMISSION_DESCRIPTIONS},
|
|
{"side": "payee", "old": RIAL, "new": SETTLEMENT_TRANSIT, "desc_none": COMMISSION_DESCRIPTIONS},
|
|
{"side": "payee", "old": REWARD, "new": SETTLEMENT_TRANSIT, "desc_none": COMMISSION_DESCRIPTIONS},
|
|
],
|
|
},
|
|
"advertising": {
|
|
"app_names": ["ad app", "advertising", "advertisement", "billboard"],
|
|
"cross_check_wallet": ADVERTISING_TRANSIT,
|
|
"legs": [
|
|
# visit-reward payout, ad-balance refund, content/tip deposit to creator
|
|
{"side": "payer", "old": REWARD, "new": ADVERTISING_TRANSIT},
|
|
# content/tip withdraw from the visitor
|
|
{"side": "payee", "old": REWARD, "new": ADVERTISING_TRANSIT},
|
|
],
|
|
},
|
|
"promotions": {
|
|
"app_names": ["promotion", "promotions"],
|
|
"cross_check_wallet": PROMOTIONS_CREDIT,
|
|
"legs": [
|
|
{"side": "payer", "old": REWARD, "new": PROMOTIONS_CREDIT},
|
|
],
|
|
},
|
|
}
|
|
|
|
|
|
class Command(BaseCommand):
|
|
help = "Repoint historical Transaction company-side accounts onto the new dedicated wallets."
|
|
|
|
def add_arguments(self, parser):
|
|
parser.add_argument("--execute", action="store_true", help="write changes (default: dry-run)")
|
|
parser.add_argument("--service", choices=sorted(SPECS), action="append",
|
|
help="limit to this service (repeatable); default: all")
|
|
for svc in SPECS:
|
|
parser.add_argument(f"--app-{svc}", dest=f"app_{svc}", metavar="UUID",
|
|
help=f"force the {svc} Application uuid instead of resolving by name")
|
|
|
|
# ── app resolution ──────────────────────────────────────────────────────
|
|
def _resolve_app(self, svc, spec, options):
|
|
override = options.get(f"app_{svc}")
|
|
if override:
|
|
try:
|
|
return Application.objects.get(pk=override)
|
|
except Application.DoesNotExist:
|
|
raise CommandError(f"[{svc}] --app-{svc}={override} is not an Application")
|
|
|
|
# 1) an existing APPLICATION account already sitting on the new dedicated wallet
|
|
acct = Account.objects.filter(
|
|
wallet_id=spec["cross_check_wallet"], owner_type=TypeChoices.APPLICATION
|
|
).first()
|
|
if acct:
|
|
try:
|
|
return Application.objects.get(pk=acct.owner_uuid)
|
|
except Application.DoesNotExist:
|
|
pass
|
|
|
|
# 2) by name
|
|
q = Q()
|
|
for name in spec["app_names"]:
|
|
q |= Q(name__iexact=name) | Q(name__icontains=name)
|
|
matches = list(Application.objects.filter(q))
|
|
if len(matches) == 1:
|
|
return matches[0]
|
|
if not matches:
|
|
raise CommandError(
|
|
f"[{svc}] could not resolve the Application (tried names {spec['app_names']}). "
|
|
f"Re-run with --app-{svc} <uuid>. Known apps: "
|
|
+ ", ".join(f'{a.name}={a.pk}' for a in Application.objects.all())
|
|
)
|
|
raise CommandError(
|
|
f"[{svc}] ambiguous Application match: "
|
|
+ ", ".join(f'{a.name}={a.pk}' for a in matches)
|
|
+ f". Re-run with --app-{svc} <uuid>."
|
|
)
|
|
|
|
# ── one leg ─────────────────────────────────────────────────────────────
|
|
def _process_leg(self, svc, app, leg, execute):
|
|
side = leg["side"]
|
|
fk = f"{side}_account"
|
|
old_acct = Account.objects.filter(
|
|
owner_uuid=app.pk, owner_type=TypeChoices.APPLICATION, wallet_id=leg["old"]
|
|
).first()
|
|
if not old_acct:
|
|
self.stdout.write(f" {svc}/{side} {leg['old']}→{leg['new']}: no old account, skip")
|
|
return
|
|
|
|
qs = Transaction.objects.filter(application=app, **{fk: old_acct})
|
|
if leg.get("desc_any"):
|
|
dq = Q()
|
|
for d in leg["desc_any"]:
|
|
dq |= Q(details__description__icontains=d)
|
|
qs = qs.filter(dq)
|
|
if leg.get("desc_none"):
|
|
for d in leg["desc_none"]:
|
|
qs = qs.exclude(details__description__icontains=d)
|
|
|
|
total = qs.count()
|
|
live_states = PAYER_LIVE_STATES if side == "payer" else PAYEE_LIVE_STATES
|
|
delta = qs.filter(state__in=live_states).aggregate(s=Sum("amount"))["s"] or 0
|
|
|
|
# sign of the balance correction
|
|
# payer (deposit): these txns had debited old_acct by `delta`
|
|
# → give it back to old, take it from new
|
|
# payee (withdraw): these txns had credited old_acct by `delta`
|
|
# → remove from old, add to new
|
|
if side == "payer":
|
|
old_delta, new_delta = +delta, -delta
|
|
else:
|
|
old_delta, new_delta = -delta, +delta
|
|
|
|
self.stdout.write(
|
|
f" {svc}/{side} {leg['old']}→{leg['new']}: "
|
|
f"{total} rows, Σ(live)={delta} "
|
|
f"balance: old {old_acct.balance}→{old_acct.balance + old_delta}, "
|
|
f"new →{new_delta:+d}"
|
|
)
|
|
|
|
if not execute or total == 0:
|
|
return
|
|
|
|
new_acct, _ = Account.objects.select_for_update().get_or_create(
|
|
owner_uuid=app.pk, owner_type=TypeChoices.APPLICATION, wallet_id=leg["new"],
|
|
defaults={"balance": 0},
|
|
)
|
|
locked_old = Account.objects.select_for_update().get(pk=old_acct.pk)
|
|
moved = qs.update(**{fk: new_acct})
|
|
Account.objects.filter(pk=locked_old.pk).update(balance=locked_old.balance + old_delta)
|
|
new_acct.refresh_from_db()
|
|
Account.objects.filter(pk=new_acct.pk).update(balance=new_acct.balance + new_delta)
|
|
self.stdout.write(self.style.SUCCESS(f" moved {moved} rows"))
|
|
|
|
# ── entrypoint ──────────────────────────────────────────────────────────
|
|
def handle(self, *args, **options):
|
|
execute = options["execute"]
|
|
services = options.get("service") or list(SPECS)
|
|
|
|
# fail early if a hard-coded wallet UUID is missing from this DB
|
|
for uid in {RIAL, REWARD, IPG_CREDIT, SETTLEMENT_TRANSIT, SETTLEMENT_COMMISSION_INCOME,
|
|
ADVERTISING_TRANSIT, PROMOTIONS_CREDIT}:
|
|
if not Wallet.objects.filter(pk=uid).exists():
|
|
raise CommandError(f"Wallet {uid} not found in this database — wrong env?")
|
|
|
|
self.stdout.write(self.style.WARNING("DRY RUN — no changes\n" if not execute else "EXECUTING\n"))
|
|
|
|
ctx = transaction.atomic() if execute else _null_ctx()
|
|
with ctx:
|
|
for svc in services:
|
|
spec = SPECS[svc]
|
|
app = self._resolve_app(svc, spec, options)
|
|
self.stdout.write(f"{svc}: Application {app.name} ({app.pk})")
|
|
for leg in spec["legs"]:
|
|
self._process_leg(svc, app, leg, execute)
|
|
self.stdout.write("")
|
|
|
|
if not execute:
|
|
self.stdout.write(self.style.WARNING("re-run with --execute to apply"))
|
|
|
|
|
|
class _null_ctx:
|
|
def __enter__(self):
|
|
return self
|
|
|
|
def __exit__(self, *a):
|
|
return False
|