refactor(payments): extract subscription purchase processing to a service

This commit is contained in:
2026-08-03 10:05:21 +07:00
parent 74e952698b
commit dee2cad540
2 changed files with 191 additions and 238 deletions

View File

@@ -1,33 +1,13 @@
# ruff: noqa: N803 # ruff: noqa: N803
import hashlib
import hmac
import logging import logging
import math
from datetime import UTC, datetime
from fastapi import Depends, Form, HTTPException from fastapi import Depends, Form, HTTPException
from fastapi.routing import APIRouter from fastapi.routing import APIRouter
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from config import cfg
from core.deps import get_db from core.deps import get_db
from db.models.orders import OrderStatus
from db.models.transactions import BalanceTransaction, BalanceTxType
from external.pally import BillStatus from external.pally import BillStatus
from external.rw import sync_subscription_by_telegram_id from services.payments import process_subscription_purchase, validate_pally_signature
from repositories import AddonsRepository
from repositories.invoices import InvoiceRepository
from repositories.orders import OrderRepository
from repositories.pricing import PricingRepository
from repositories.users import UserRepository
from schemas.invoices import InvoiceStatus
from services.plans import get_pricing_model
from services.subscriptions import (
apply_order_now,
deduct_order_balance,
queue_order_for_later,
should_apply_immediately,
)
router = APIRouter(prefix="/payments/pally") router = APIRouter(prefix="/payments/pally")
@@ -35,7 +15,7 @@ logger = logging.getLogger(__name__)
@router.post("/result") @router.post("/result")
async def pally_callback( # noqa: PLR0911, PLR0912, PLR0915 async def pally_callback(
*, *,
InvId: str = Form(...), InvId: str = Form(...),
OutSum: str = Form(...), OutSum: str = Form(...),
@@ -45,7 +25,6 @@ async def pally_callback( # noqa: PLR0911, PLR0912, PLR0915
CurrencyIn: str = Form(...), CurrencyIn: str = Form(...),
custom: str | None = Form(None), custom: str | None = Form(None),
SignatureValue: str = Form(...), SignatureValue: str = Form(...),
# Optional fields for additional information
AccountType: str | None = Form(None), AccountType: str | None = Form(None),
AccountNumber: str | None = Form(None), AccountNumber: str | None = Form(None),
BalanceAmount: str | None = Form(None), BalanceAmount: str | None = Form(None),
@@ -58,14 +37,6 @@ async def pally_callback( # noqa: PLR0911, PLR0912, PLR0915
ErrorMessage: str | None = Form(None), ErrorMessage: str | None = Form(None),
session: AsyncSession = Depends(get_db), session: AsyncSession = Depends(get_db),
): ):
invoice_repo = InvoiceRepository(session)
orders_repo = OrderRepository(session)
users_repo = UserRepository(session)
pricing_repo = PricingRepository(session)
invoice_id_str = InvId
invoice_id: int | None = None
invoice_creator_id: int | None = None
logger.info( logger.info(
"Pally webhook received - InvId: %s, OutSum: %s, Commission: %s, TrsId: %s, Status: %s, " "Pally webhook received - InvId: %s, OutSum: %s, Commission: %s, TrsId: %s, Status: %s, "
"CurrencyIn: %s, custom: %s, BalanceAmount: %s, SignatureValue: %s", "CurrencyIn: %s, custom: %s, BalanceAmount: %s, SignatureValue: %s",
@@ -89,29 +60,18 @@ async def pally_callback( # noqa: PLR0911, PLR0912, PLR0915
TrsId, TrsId,
) )
# Validate signature if not validate_pally_signature(OutSum, InvId, SignatureValue):
raw_string = f"{OutSum}:{InvId}:{cfg.pally_token}"
expected_signature = hashlib.md5(raw_string.encode("utf-8")).hexdigest().upper()
logger.debug("Signature validation for TrsId %s", TrsId)
if not hmac.compare_digest(SignatureValue, expected_signature):
logger.critical( logger.critical(
"SECURITY ALERT: Invalid signature for TrsId %s - Expected: %s, Received: %s", "SECURITY ALERT: Invalid signature for TrsId %s",
TrsId, TrsId,
expected_signature,
SignatureValue,
) )
raise HTTPException(403, detail="Invalid signature.") raise HTTPException(403, detail="Invalid signature.")
# Only process successful payments
if Status != BillStatus.SUCCESS: if Status != BillStatus.SUCCESS:
logger.info("Bill %s skipped: status=%s", TrsId, Status) logger.info("Bill %s skipped: status=%s", TrsId, Status)
return "OK" return "OK"
logger.info("Processing successfully paid bill %s", TrsId) invoice_id_str = InvId
# Validate invoice ID
if not invoice_id_str or not invoice_id_str.isdigit(): if not invoice_id_str or not invoice_id_str.isdigit():
logger.critical( logger.critical(
"Invalid or non-numeric bill ID in InvId field for TrsId %s: '%s'", "Invalid or non-numeric bill ID in InvId field for TrsId %s: '%s'",
@@ -120,199 +80,13 @@ async def pally_callback( # noqa: PLR0911, PLR0912, PLR0915
) )
return "OK" return "OK"
# Find bill in database amount = int(float(BalanceAmount)) if BalanceAmount is not None else int(float(OutSum))
invoice = await invoice_repo.get_by_id(int(invoice_id_str))
if not invoice: await process_subscription_purchase(
logger.critical("Bill %s not found in database for TrsId %s", invoice_id_str, TrsId) session,
return "OK" invoice_id=int(invoice_id_str),
trs_id=TrsId,
invoice_id = invoice.id amount=amount,
invoice_creator_id = invoice.creator_id )
# Check if already processed
if invoice.status != InvoiceStatus.ACTIVE:
logger.warning(
"Bill %s (TrsId: %s) is already processed with status: %s",
invoice.id,
TrsId,
invoice.status,
)
return "OK"
try:
# Handle fee scenarios: use BalanceAmount if available (net amount after fees),
# otherwise use OutSum (gross amount paid by customer)
if BalanceAmount is not None:
# Customer pays fees - BalanceAmount is the net amount credited to merchant
credited_amount = int(float(BalanceAmount))
gross_amount = int(float(OutSum))
logger.info(
"Customer-pays-fees payment: bill_id=%s, expected=%s, gross_paid=%s, net_credited=%s",
invoice.id,
invoice.amount,
gross_amount,
credited_amount,
)
# Validate that the net credited amount matches our bill amount
if invoice.amount != credited_amount:
logger.error(
"Net amount mismatch for bill %s (TrsId: %s) - Expected: %s, Net credited: %s, Gross paid: %s",
invoice.id,
TrsId,
invoice.amount,
credited_amount,
gross_amount,
)
return "OK"
amount = credited_amount # Credit the net amount (without fees)
else:
# Standard payment - OutSum should match bill amount exactly
amount = int(float(OutSum))
logger.info(
"Standard payment: bill_id=%s, expected=%s, received=%s",
invoice.id,
invoice.amount,
amount,
)
if invoice.amount != amount:
logger.error(
"Amount mismatch for bill %s (TrsId: %s) - Expected: %s, Received: %s",
invoice.id,
TrsId,
invoice.amount,
amount,
)
return "OK"
if invoice.order_id is None:
logger.critical("Invoice %s has no linked order for TrsId %s", invoice.id, TrsId)
return "OK"
order = await orders_repo.get_by_id(invoice.order_id)
if order is None:
logger.critical(
"Order not found for bill %s (TrsId: %s, user_id=%s, order_id=%s)",
invoice.id,
TrsId,
invoice.creator_id,
invoice.order_id,
)
return "OK"
if order.user_id != invoice.creator_id:
logger.critical(
"Order %s does not belong to invoice creator %s for TrsId %s",
order.id,
invoice.creator_id,
TrsId,
)
return "OK"
user = await users_repo.get_user_by_id(invoice.creator_id)
if user is None:
logger.critical("User %s not found for TrsId %s", invoice.creator_id, TrsId)
return "OK"
subscription = user.subscription
now = datetime.now(UTC)
pricing = await get_pricing_model(AddonsRepository(session), pricing_repo)
logger.info(
"Processing payment: bill_id=%s, order_id=%s, user_id=%s, amount=%s",
invoice.id,
order.id,
invoice.creator_id,
amount,
)
order.status = OrderStatus.PAID
invoice.status = InvoiceStatus.PAID
if order.balance_amount > 0:
await deduct_order_balance(
session,
user=user,
order=order,
description=f"order {order.id} partial payment from balance",
)
if should_apply_immediately(
subscription=subscription,
order=order,
pricing=pricing,
now=now,
):
applied_subscription = await apply_order_now(
session, user=user, order=order, pricing=pricing, now=now
)
await sync_subscription_by_telegram_id(
expires_at=applied_subscription.expires_at,
devices=applied_subscription.devices,
telegram_id=user.telegram_id,
username=user.username,
)
else:
if subscription is None:
logger.critical(
"Cannot queue order %s without subscription for user %s", order.id, user.id
)
return "OK"
await queue_order_for_later(order=order, subscription=subscription, now=now)
# Process referral bonus
referal_id = user.referal_id
if referal_id is not None:
referal_amount = math.floor(amount * (cfg.referal_bonus / 100))
if referal_amount > 0:
referal_user = await session.get(type(user), referal_id)
if referal_user is not None:
session.add(
BalanceTransaction(
user_id=referal_id,
amount=referal_amount,
tx_type=BalanceTxType.REFERRAL_BONUS,
balance_before=referal_user.balance,
balance_after=referal_user.balance + referal_amount,
description=f"referral reward for user {invoice.creator_id} (TrsId: {TrsId})",
)
)
referal_user.balance += referal_amount
logger.info(
"Referral bonus processed: referrer_id=%s, amount=%s",
referal_id,
referal_amount,
)
else:
logger.warning(
"Referrer %s not found for user_id=%s while processing TrsId %s",
referal_id,
invoice.creator_id,
TrsId,
)
await session.commit()
logger.info("Payment processing completed successfully for TrsId %s", TrsId)
except Exception as e:
await session.rollback()
logger.exception(
"CRITICAL ERROR processing payment for TrsId %s, bill_id %s, user_id %s: %s",
TrsId,
invoice_id,
invoice_creator_id,
str(e),
)
raise HTTPException(500, detail="Payment processing failed.") from e
logger.info("Bill %s marked as PAID for TrsId %s", invoice_id, TrsId)
return "OK" return "OK"

179
services/payments.py Normal file
View File

@@ -0,0 +1,179 @@
import hashlib
import hmac
import logging
import math
from datetime import UTC, datetime
from config import cfg
from db.models.orders import OrderStatus
from db.models.transactions import BalanceTransaction, BalanceTxType
from external.rw import sync_subscription_by_telegram_id
from repositories import AddonsRepository
from repositories.invoices import InvoiceRepository
from repositories.orders import OrderRepository
from repositories.pricing import PricingRepository
from repositories.users import UserRepository
from schemas.invoices import InvoiceStatus
from services.plans import get_pricing_model
from services.subscriptions import (
apply_order_now,
deduct_order_balance,
queue_order_for_later,
should_apply_immediately,
)
logger = logging.getLogger(__name__)
def validate_pally_signature(out_sum: str, inv_id: str, signature_value: str) -> bool:
raw_string = f"{out_sum}:{inv_id}:{cfg.pally_token}"
expected_signature = hashlib.md5(raw_string.encode("utf-8")).hexdigest().upper()
return hmac.compare_digest(signature_value, expected_signature)
async def process_subscription_purchase( # noqa: PLR0911, PLR0912
session,
*,
invoice_id: int,
trs_id: str,
amount: int,
) -> None:
invoice_repo = InvoiceRepository(session)
orders_repo = OrderRepository(session)
users_repo = UserRepository(session)
pricing_repo = PricingRepository(session)
invoice = await invoice_repo.get_by_id(invoice_id)
if not invoice:
logger.critical("Bill %s not found in database for TrsId %s", invoice_id, trs_id)
return
if invoice.status != InvoiceStatus.ACTIVE:
logger.warning(
"Bill %s (TrsId: %s) is already processed with status: %s",
invoice.id,
trs_id,
invoice.status,
)
return
if invoice.amount != amount:
logger.error(
"Amount mismatch for bill %s (TrsId: %s) - Expected: %s, Received: %s",
invoice.id,
trs_id,
invoice.amount,
amount,
)
return
if invoice.order_id is None:
logger.critical("Invoice %s has no linked order for TrsId %s", invoice.id, trs_id)
return
order = await orders_repo.get_by_id(invoice.order_id)
if order is None:
logger.critical(
"Order not found for bill %s (TrsId: %s, user_id=%s, order_id=%s)",
invoice.id,
trs_id,
invoice.creator_id,
invoice.order_id,
)
return
if order.user_id != invoice.creator_id:
logger.critical(
"Order %s does not belong to invoice creator %s for TrsId %s",
order.id,
invoice.creator_id,
trs_id,
)
return
user = await users_repo.get_user_by_id(invoice.creator_id)
if user is None:
logger.critical("User %s not found for TrsId %s", invoice.creator_id, trs_id)
return
subscription = user.subscription
now = datetime.now(UTC)
pricing = await get_pricing_model(AddonsRepository(session), pricing_repo)
logger.info(
"Processing payment: bill_id=%s, order_id=%s, user_id=%s, amount=%s",
invoice.id,
order.id,
invoice.creator_id,
amount,
)
order.status = OrderStatus.PAID
invoice.status = InvoiceStatus.PAID
if order.balance_amount > 0:
await deduct_order_balance(
session,
user=user,
order=order,
description=f"order {order.id} partial payment from balance",
)
if should_apply_immediately(
subscription=subscription,
order=order,
pricing=pricing,
now=now,
):
applied_subscription = await apply_order_now(
session, user=user, order=order, pricing=pricing, now=now
)
await sync_subscription_by_telegram_id(
expires_at=applied_subscription.expires_at,
devices=applied_subscription.devices,
telegram_id=user.telegram_id,
username=user.username,
)
else:
if subscription is None:
logger.critical(
"Cannot queue order %s without subscription for user %s", order.id, user.id
)
return
await queue_order_for_later(order=order, subscription=subscription, now=now)
referal_id = user.referal_id
if referal_id is not None:
referal_amount = math.floor(amount * (cfg.referal_bonus / 100))
if referal_amount > 0:
referal_user = await session.get(type(user), referal_id)
if referal_user is not None:
session.add(
BalanceTransaction(
user_id=referal_id,
amount=referal_amount,
tx_type=BalanceTxType.REFERRAL_BONUS,
balance_before=referal_user.balance,
balance_after=referal_user.balance + referal_amount,
description=f"referral reward for user {invoice.creator_id} (TrsId: {trs_id})",
)
)
referal_user.balance += referal_amount
logger.info(
"Referral bonus processed: referrer_id=%s, amount=%s",
referal_id,
referal_amount,
)
else:
logger.warning(
"Referrer %s not found for user_id=%s while processing TrsId %s",
referal_id,
invoice.creator_id,
trs_id,
)
await session.commit()
logger.info("Bill %s marked as PAID for TrsId %s", invoice_id, trs_id)