refactor: split backend domains and add API behavior coverage

This commit is contained in:
3252a8
2026-05-13 18:38:10 +03:00
parent a96007d48a
commit 7401649841
103 changed files with 10558 additions and 9832 deletions
+46 -2101
View File
File diff suppressed because it is too large Load Diff
+1
View File
@@ -0,0 +1 @@
"""Domain modules for the admin Mini App API."""
+65
View File
@@ -0,0 +1,65 @@
# ruff: noqa: F401,F403,F405,I001
"""HTTP API powering the admin section of the subscription Mini App.
All routes require an authenticated webapp session (cookie or Bearer
token) AND the resolved Telegram user id must appear in
``settings.ADMIN_IDS``. Authorization is enforced via the
``_require_admin_user_id`` helper, never trusted from the client.
"""
from __future__ import annotations
import csv
import io
import json
import logging
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
from urllib.parse import parse_qsl, urlsplit, urlunsplit
from aiohttp import web
from pydantic import ValidationError
from sqlalchemy import Float, and_, case, cast, or_, select
from sqlalchemy import func as sa_func
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import sessionmaker
from bot.app.web.admin_settings_manifest import (
manifest_payload,
)
from bot.services.referral_service import ReferralService
from bot.services.settings_override_service import (
current_value,
update_overrides,
)
from bot.utils import MessageContent, send_message_via_queue
from bot.utils.message_queue import get_queue_manager
from config.settings import Settings
from config.tariffs_config import TariffsConfig
from db.dal import (
ad_dal,
app_settings_dal,
message_log_dal,
panel_sync_dal,
payment_dal,
promo_code_dal,
subscription_dal,
user_dal,
)
from db.models import (
AdCampaign,
MessageLog,
Payment,
PromoCode,
Subscription,
User,
UserTelegramAvatar,
)
logger = logging.getLogger(__name__)
# ─── Auth ──────────────────────────────────────────────────────────
__all__ = [name for name in globals() if not name.startswith("__")]
+69
View File
@@ -0,0 +1,69 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_ads_list_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
campaigns = await ad_dal.list_campaigns(session)
totals = await ad_dal.get_totals(session)
results = []
for campaign in campaigns:
try:
stats = await ad_dal.get_campaign_stats(session, campaign.ad_campaign_id)
except Exception:
stats = {}
results.append(_serialize_ad(campaign, stats))
return _ok({"campaigns": results, "totals": totals})
async def admin_ad_create_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
payload = await _read_json(request)
source = str(payload.get("source") or "").strip()
start_param = str(payload.get("start_param") or "").strip()
cost = float(payload.get("cost") or 0.0)
if not source or not start_param:
return _error(400, "invalid_payload")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
existing = await ad_dal.get_campaign_by_start_param(session, start_param)
if existing:
return _error(409, "duplicate_start_param")
campaign = await ad_dal.create_campaign(
session,
source=source,
start_param=start_param,
cost=cost,
)
await session.commit()
await session.refresh(campaign)
return _ok({"campaign": _serialize_ad(campaign)})
async def admin_ad_toggle_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
campaign_id = int(request.match_info["campaign_id"])
payload = await _read_json(request)
is_active = bool(payload.get("is_active", True))
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
ok = await ad_dal.toggle_campaign_active(session, campaign_id, is_active)
if not ok:
return _error(404, "not_found")
await session.commit()
return _ok({})
async def admin_ad_delete_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
campaign_id = int(request.match_info["campaign_id"])
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
ok = await ad_dal.delete_campaign(session, campaign_id)
if not ok:
return _error(404, "not_found")
await session.commit()
return _ok({})
+58
View File
@@ -0,0 +1,58 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
def _require_admin_user_id(request: web.Request) -> int:
"""Return the authenticated user id, or raise 401/403 for non-admins."""
from bot.app.web.session import extract_authenticated_user_id
settings: Settings = request.app["settings"]
user_id = extract_authenticated_user_id(request)
if not user_id:
raise web.HTTPUnauthorized(
text=json.dumps({"ok": False, "error": "unauthorized"}),
content_type="application/json",
)
admin_ids = settings.ADMIN_IDS or []
db_user_telegram_id = request.get("admin_telegram_id")
if db_user_telegram_id is None:
raise web.HTTPForbidden(
text=json.dumps({"ok": False, "error": "forbidden"}),
content_type="application/json",
)
if int(db_user_telegram_id) not in {int(x) for x in admin_ids}:
raise web.HTTPForbidden(
text=json.dumps({"ok": False, "error": "forbidden"}),
content_type="application/json",
)
return int(user_id)
@web.middleware
async def admin_auth_middleware(request: web.Request, handler):
"""Resolve the Telegram id of the current user and stash it on the request.
Doing this once per request lets every admin route call
``_require_admin_user_id`` without re-querying the DB.
"""
if not request.path.startswith("/api/admin"):
return await handler(request)
from bot.app.web.session import extract_authenticated_user_id
user_id = extract_authenticated_user_id(request)
if user_id:
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if db_user and db_user.telegram_id:
request["admin_telegram_id"] = int(db_user.telegram_id)
elif db_user:
# No telegram_id yet (email-only user) — can't be an admin
request["admin_telegram_id"] = None
return await handler(request)
+54
View File
@@ -0,0 +1,54 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_broadcast_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
payload = await _read_json(request)
text = str(payload.get("text") or "").strip()
target = str(payload.get("target") or "all").strip().lower()
if not text:
return _error(400, "empty_text")
if target not in {"all", "active", "inactive"}:
target = "all"
queue_manager = get_queue_manager()
if not queue_manager:
return _error(503, "queue_unavailable")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
if target == "active":
user_ids = await user_dal.get_user_ids_with_active_subscription(session)
elif target == "inactive":
user_ids = await user_dal.get_user_ids_without_active_subscription(session)
else:
user_ids = await user_dal.get_all_active_user_ids_for_broadcast(session)
sent = 0
failed = 0
for uid in user_ids:
try:
await send_message_via_queue(
queue_manager,
int(uid),
MessageContent(content_type="text", text=text),
parse_mode="HTML",
disable_web_page_preview=True,
)
sent += 1
except Exception as exc:
failed += 1
logger.debug("Broadcast queue failed for %s: %s", uid, exc)
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_broadcast_webapp",
"content": f"target={target} sent={sent} failed={failed} text={text[:120]}",
"is_admin_event": True,
},
)
return _ok({"queued": sent, "failed": failed, "target": target})
+367
View File
@@ -0,0 +1,367 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
def _ok(payload: Dict[str, Any], **extra) -> web.Response:
body = {"ok": True, **payload, **extra}
return web.json_response(body)
def _error(status: int, code: str, message: str = "") -> web.Response:
return web.json_response(
{"ok": False, "error": code, "message": message or code},
status=status,
)
async def _read_json(request: web.Request) -> Dict[str, Any]:
try:
data = await request.json()
return data if isinstance(data, dict) else {}
except Exception:
return {}
def _serialize_user(user: User) -> Dict[str, Any]:
return {
"user_id": int(user.user_id),
"telegram_id": int(user.telegram_id) if user.telegram_id else None,
"telegram_photo_url": user.telegram_photo_url,
"username": user.username,
"first_name": user.first_name,
"last_name": user.last_name,
"email": user.email,
"language_code": user.language_code,
"is_banned": bool(user.is_banned),
"registration_date": user.registration_date.isoformat() if user.registration_date else None,
"panel_user_uuid": user.panel_user_uuid,
"referral_code": user.referral_code,
"referred_by_id": int(user.referred_by_id) if user.referred_by_id else None,
}
def _premium_limit_bytes_from_subscription(sub: Subscription) -> int:
premium_bonus_bytes = int(getattr(sub, "premium_bonus_bytes", 0) or 0)
return (
int(sub.premium_baseline_bytes or 0)
+ int(sub.premium_topup_balance_bytes or 0)
+ int(getattr(sub, "premium_topup_used_bytes", 0) or 0)
+ premium_bonus_bytes
)
def _premium_traffic_list_payload(sub: Optional[Subscription]) -> Dict[str, Any]:
"""Premium traffic column when subscription has a finite premium quota (bytes > 0).
Note: ``Subscription.premium_is_limited`` in the DB means *quota exhausted* for panel
routing, not 'tariff includes premium traffic' — do not use it here.
"""
if sub is None:
return {"state": "none"}
if bool(getattr(sub, "premium_unlimited_override", False)):
return {
"state": "unlimited",
"unlimited": True,
"used_bytes": int(sub.premium_used_bytes or 0),
"limit_bytes": None,
"percent": None,
}
limit_bytes = _premium_limit_bytes_from_subscription(sub)
if limit_bytes <= 0:
return {"state": "none"}
used_bytes = int(sub.premium_used_bytes or 0)
ratio = float(used_bytes) / float(limit_bytes) if limit_bytes else 0.0
pct = int(max(0, min(100, round(ratio * 100))))
if ratio >= 1.0:
state = "critical"
elif ratio >= 0.85:
state = "warn"
else:
state = "good"
return {
"state": state,
"unlimited": False,
"used_bytes": used_bytes,
"limit_bytes": limit_bytes,
"percent": pct,
}
def _serialize_subscription(sub: Subscription) -> Dict[str, Any]:
premium_bonus_bytes = int(getattr(sub, "premium_bonus_bytes", 0) or 0)
regular_bonus_bytes = int(getattr(sub, "regular_bonus_bytes", 0) or 0)
regular_unlimited_override = bool(getattr(sub, "regular_unlimited_override", False))
premium_unlimited_override = bool(getattr(sub, "premium_unlimited_override", False))
premium_limit_bytes = _premium_limit_bytes_from_subscription(sub)
return {
"subscription_id": int(sub.subscription_id),
"panel_user_uuid": sub.panel_user_uuid,
"panel_subscription_uuid": sub.panel_subscription_uuid,
"start_date": sub.start_date.isoformat() if sub.start_date else None,
"end_date": sub.end_date.isoformat() if sub.end_date else None,
"duration_months": sub.duration_months,
"is_active": bool(sub.is_active),
"status_from_panel": sub.status_from_panel,
"traffic_limit_bytes": sub.traffic_limit_bytes,
"traffic_used_bytes": sub.traffic_used_bytes,
"tier_baseline_bytes": sub.tier_baseline_bytes,
"topup_balance_bytes": sub.topup_balance_bytes,
"premium_used_bytes": sub.premium_used_bytes,
"premium_limit_bytes": premium_limit_bytes,
"premium_baseline_bytes": sub.premium_baseline_bytes,
"premium_topup_balance_bytes": sub.premium_topup_balance_bytes,
"premium_topup_used_bytes": getattr(sub, "premium_topup_used_bytes", 0),
"premium_bonus_bytes": premium_bonus_bytes,
"regular_bonus_bytes": regular_bonus_bytes,
"regular_unlimited_override": regular_unlimited_override,
"premium_unlimited_override": premium_unlimited_override,
"premium_is_limited": bool(sub.premium_is_limited),
"tariff_key": sub.tariff_key,
"auto_renew_enabled": bool(sub.auto_renew_enabled),
"provider": sub.provider,
"is_throttled": bool(sub.is_throttled),
}
def _payment_traffic_gb_split(payment: Payment) -> Tuple[Optional[float], Optional[float]]:
"""For traffic purchases: ``(regular_gb, premium_gb)``. Other payments → (None, None)."""
if payment.purchased_gb is None:
return None, None
try:
gb = float(payment.purchased_gb)
except (TypeError, ValueError):
return None, None
sm = (payment.sale_mode or "").strip()
if not sm:
return None, None
base = sm.split("@", 1)[0].split("|", 1)[0].lower()
if base == "premium_topup":
return None, gb
if base in {"traffic", "traffic_package", "topup"}:
return gb, None
return None, None
def _payment_user_display_label(loaded_user: Any, payment_user_id: int) -> str:
"""Human-facing name for payments tables: TG profile name, else email, else user id."""
if loaded_user is None:
return str(payment_user_id)
tid = getattr(loaded_user, "telegram_id", None)
if tid is not None:
fn = (getattr(loaded_user, "first_name", None) or "").strip()
ln = (getattr(loaded_user, "last_name", None) or "").strip()
full = f"{fn} {ln}".strip()
if full:
return full
un = (getattr(loaded_user, "username", None) or "").strip()
if un:
return un if un.startswith("@") else f"@{un}"
return str(payment_user_id)
email = (getattr(loaded_user, "email", None) or "").strip()
if email:
return email
return str(payment_user_id)
def _serialize_payment(payment: Payment) -> Dict[str, Any]:
# Avoid lazy-loading `payment.user` outside an active SQLAlchemy session.
# Some admin routes serialize payments after the session scope is closed.
telegram_id = None
loaded_user = payment.__dict__.get("user")
user_label = _payment_user_display_label(loaded_user, int(payment.user_id))
if loaded_user is not None:
tid = getattr(loaded_user, "telegram_id", None)
if tid is not None:
try:
telegram_id = int(tid)
except (TypeError, ValueError):
telegram_id = None
reg_gb, prem_gb = _payment_traffic_gb_split(payment)
return {
"payment_id": int(payment.payment_id),
"user_id": int(payment.user_id),
"user_label": user_label,
"telegram_id": telegram_id,
"traffic_regular_gb": reg_gb,
"traffic_premium_gb": prem_gb,
"provider": payment.provider,
"provider_payment_id": payment.provider_payment_id,
"amount": float(payment.amount),
"currency": payment.currency,
"status": payment.status,
"description": payment.description,
"subscription_duration_months": payment.subscription_duration_months,
"sale_mode": payment.sale_mode,
"tariff_key": payment.tariff_key,
"purchased_gb": payment.purchased_gb,
"purchased_hwid_devices": payment.purchased_hwid_devices,
"created_at": payment.created_at.isoformat() if payment.created_at else None,
}
def _serialize_promo(promo: PromoCode) -> Dict[str, Any]:
return {
"id": int(promo.promo_code_id),
"code": promo.code,
"bonus_days": int(promo.bonus_days),
"max_activations": int(promo.max_activations),
"current_activations": int(promo.current_activations or 0),
"is_active": bool(promo.is_active),
"valid_until": promo.valid_until.isoformat() if promo.valid_until else None,
"created_at": promo.created_at.isoformat() if promo.created_at else None,
"created_by_admin_id": int(promo.created_by_admin_id)
if promo.created_by_admin_id
else None,
}
def _serialize_ad(campaign: AdCampaign, totals: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
return {
"id": int(campaign.ad_campaign_id),
"source": campaign.source,
"start_param": campaign.start_param,
"cost": float(campaign.cost or 0),
"is_active": bool(campaign.is_active),
"created_at": campaign.created_at.isoformat() if campaign.created_at else None,
"stats": totals or {},
}
def _serialize_log(entry: MessageLog) -> Dict[str, Any]:
return {
"log_id": int(entry.log_id),
"user_id": int(entry.user_id) if entry.user_id else None,
"telegram_username": entry.telegram_username,
"telegram_first_name": entry.telegram_first_name,
"event_type": entry.event_type,
"content": entry.content,
"is_admin_event": bool(entry.is_admin_event),
"target_user_id": int(entry.target_user_id) if entry.target_user_id else None,
"timestamp": entry.timestamp.isoformat() if entry.timestamp else None,
}
def _tariffs_config_path(settings: Settings) -> Path:
return Path(settings.TARIFFS_CONFIG_PATH).expanduser()
def _tariffs_config_payload(config: TariffsConfig) -> Dict[str, Any]:
return config.model_dump(mode="json", exclude_none=True)
def _write_tariffs_config_file(path: Path, config: TariffsConfig) -> None:
data = _tariffs_config_payload(config)
path.parent.mkdir(parents=True, exist_ok=True)
tmp_path = path.with_suffix(f"{path.suffix}.tmp")
payload = json.dumps(data, ensure_ascii=False, indent=2) + "\n"
try:
tmp_path.write_text(payload, encoding="utf-8")
tmp_path.replace(path)
except PermissionError:
# A docker-compose single-file bind mount can make /app/config
# unwritable while the mounted tariffs.json itself is writable.
# Fall back to updating the existing file in-place.
if tmp_path.exists():
try:
tmp_path.unlink()
except OSError:
pass
path.write_text(payload, encoding="utf-8")
def _panel_node_uuid_key(node: Dict[str, Any]) -> str:
uid = node.get("nodeUuid") or node.get("node_uuid") or node.get("uuid") or node.get("id")
return str(uid).strip().lower() if uid else ""
def _panel_node_users_online(node: Dict[str, Any]) -> Optional[int]:
uo = node.get("usersOnline")
if uo is None:
uo = node.get("users_online")
if uo is None:
uo = node.get("onlineUsers") or node.get("online_users")
if uo is None:
mg = node.get("metricGroups")
if isinstance(mg, dict):
uo = mg.get("onlineUsers") or mg.get("online_users")
if uo is None:
return None
try:
return int(uo)
except (TypeError, ValueError):
return None
def _panel_nodes_online_by_uuid(nodes_payload: Any) -> Dict[str, int]:
"""Build node_uuid(lower) -> usersOnline from GET /system/stats/nodes payload."""
out: Dict[str, int] = {}
raw_list: Optional[List[Any]] = None
if isinstance(nodes_payload, list):
raw_list = nodes_payload
elif isinstance(nodes_payload, dict):
raw_list = nodes_payload.get("nodes")
if raw_list is None:
raw_list = nodes_payload.get("items") or nodes_payload.get("data")
if not isinstance(raw_list, list):
return out
for n in raw_list:
if not isinstance(n, dict):
continue
key = _panel_node_uuid_key(n)
if not key:
continue
online = _panel_node_users_online(n)
if online is not None:
out[key] = online
return out
def _enrich_bandwidth_nodes_with_online(
bw: Any,
online_by_uuid: Dict[str, int],
online_by_name: Optional[Dict[str, int]] = None,
) -> None:
"""Attach usersOnline to topNodes/series (UUID and optional node name)."""
if not isinstance(bw, dict):
return
if not online_by_uuid and not online_by_name:
return
for key in ("topNodes", "series"):
arr = bw.get(key)
if not isinstance(arr, list):
continue
for item in arr:
if not isinstance(item, dict):
continue
if item.get("usersOnline") is not None:
continue
uid = item.get("uuid") or item.get("nodeUuid") or item.get("node_uuid")
if uid and online_by_uuid:
hit = online_by_uuid.get(str(uid).strip().lower())
if hit is not None:
item["usersOnline"] = hit
continue
if online_by_name:
nm = item.get("name")
if nm and isinstance(nm, str):
hitn = online_by_name.get(nm.strip().lower())
if hitn is not None:
item["usersOnline"] = hitn
def _build_admin_webapp_referral_link(
base_url: Optional[str], referral_code: Optional[str]
) -> Optional[str]:
"""Mirror of ``subscription_webapp._build_webapp_referral_link``.
Kept local to avoid a cross-module import cycle (subscription_webapp
imports admin_api).
"""
if not base_url or not referral_code:
return None
parts = urlsplit(base_url)
query = dict(parse_qsl(parts.query, keep_blank_values=True))
query["ref"] = f"u{referral_code}"
new_query = "&".join(f"{k}={v}" for k, v in query.items())
return urlunsplit((parts.scheme, parts.netloc, parts.path, new_query, parts.fragment))
+36
View File
@@ -0,0 +1,36 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_logs_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
page = max(0, int(request.query.get("page", 0) or 0))
page_size = min(200, max(1, int(request.query.get("page_size", 50) or 50)))
user_filter = request.query.get("user_id")
async with async_session_factory() as session:
if user_filter:
try:
user_id = int(user_filter)
except (TypeError, ValueError):
return _error(400, "invalid_user_id")
entries = await message_log_dal.get_user_message_logs(
session, user_id, page_size, page * page_size
)
total = await message_log_dal.count_user_message_logs(session, user_id)
else:
entries = await message_log_dal.get_all_message_logs(
session, page_size, page * page_size
)
total = await message_log_dal.count_all_message_logs(session)
return _ok(
{
"logs": [_serialize_log(entry) for entry in entries],
"page": page,
"page_size": page_size,
"total": int(total or 0),
}
)
+35
View File
@@ -0,0 +1,35 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_panel_internal_squads_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
panel_service = request.app.get("panel_service")
if panel_service is None:
return _error(503, "panel_unavailable", "Panel service unavailable")
try:
squads = await panel_service.get_internal_squads()
except Exception as exc:
logger.exception("Failed to load internal squads from panel")
return _error(502, "panel_request_failed", str(exc))
if squads is None:
return _error(502, "panel_request_failed", "Unable to load internal squads")
items = []
for squad in squads:
if not isinstance(squad, dict):
continue
uuid = squad.get("uuid") or squad.get("id")
if not uuid:
continue
items.append(
{
"uuid": str(uuid),
"name": squad.get("name") or squad.get("title") or str(uuid),
"members_count": squad.get("membersCount")
or squad.get("usersCount")
or squad.get("members_count"),
"active_inbounds_count": squad.get("activeInboundsCount")
or squad.get("active_inbounds_count"),
}
)
return _ok({"squads": items})
+95
View File
@@ -0,0 +1,95 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_payments_list_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
page = max(0, int(request.query.get("page", 0) or 0))
page_size = min(100, max(1, int(request.query.get("page_size", 25) or 25)))
async with async_session_factory() as session:
from sqlalchemy.orm import selectinload
stmt = (
select(Payment)
.options(selectinload(Payment.user))
.order_by(Payment.created_at.desc())
.offset(page * page_size)
.limit(page_size)
)
rows = (await session.execute(stmt)).scalars().all()
total = await payment_dal.get_payments_count(session)
return _ok(
{
"payments": [_serialize_payment(p) for p in rows],
"page": page,
"page_size": page_size,
"total": int(total or 0),
}
)
async def admin_payments_export_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
from sqlalchemy.orm import selectinload
stmt = (
select(Payment)
.options(selectinload(Payment.user))
.order_by(Payment.created_at.desc())
.limit(10000)
)
rows = (await session.execute(stmt)).scalars().all()
buffer = io.StringIO()
writer = csv.writer(buffer)
writer.writerow(
[
"payment_id",
"user_id",
"user_label",
"provider",
"provider_payment_id",
"amount",
"currency",
"status",
"description",
"duration_months",
"sale_mode",
"tariff_key",
"created_at",
]
)
for p in rows:
label = _payment_user_display_label(p.user, int(p.user_id)) if p.user else str(p.user_id)
writer.writerow(
[
p.payment_id,
p.user_id,
label,
p.provider,
p.provider_payment_id or "",
p.amount,
p.currency,
p.status,
p.description or "",
p.subscription_duration_months or "",
p.sale_mode or "",
p.tariff_key or "",
p.created_at.isoformat() if p.created_at else "",
]
)
response = web.Response(
body=buffer.getvalue().encode("utf-8-sig"),
content_type="text/csv",
charset="utf-8",
)
response.headers["Content-Disposition"] = 'attachment; filename="payments.csv"'
return response
+96
View File
@@ -0,0 +1,96 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_promos_list_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
page = max(0, int(request.query.get("page", 0) or 0))
page_size = min(100, max(1, int(request.query.get("page_size", 25) or 25)))
async with async_session_factory() as session:
promos = await promo_code_dal.get_all_promo_codes_with_details(
session, limit=page_size, offset=page * page_size
)
total = await promo_code_dal.get_promo_codes_count(session)
return _ok(
{
"promos": [_serialize_promo(p) for p in promos],
"page": page,
"page_size": page_size,
"total": int(total or 0),
}
)
async def admin_promo_create_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
payload = await _read_json(request)
code = str(payload.get("code") or "").strip().upper()
bonus_days = int(payload.get("bonus_days") or 0)
max_activations = int(payload.get("max_activations") or 0)
valid_days = payload.get("valid_days")
if not code or bonus_days <= 0 or max_activations <= 0:
return _error(400, "invalid_payload")
valid_until = None
if valid_days:
try:
valid_until = datetime.now(timezone.utc) + timedelta(days=int(valid_days))
except (TypeError, ValueError):
return _error(400, "invalid_valid_days")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
existing = await promo_code_dal.get_promo_code_by_code(session, code)
if existing:
return _error(409, "duplicate_code")
promo = await promo_code_dal.create_promo_code(
session,
{
"code": code,
"bonus_days": bonus_days,
"max_activations": max_activations,
"valid_until": valid_until,
"created_by_admin_id": actor_id,
"is_active": True,
},
)
await session.commit()
return _ok({"promo": _serialize_promo(promo)})
async def admin_promo_update_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
promo_id = int(request.match_info["promo_id"])
payload = await _read_json(request)
update_data: Dict[str, Any] = {}
if "is_active" in payload:
update_data["is_active"] = bool(payload["is_active"])
if "bonus_days" in payload and payload["bonus_days"] is not None:
update_data["bonus_days"] = int(payload["bonus_days"])
if "max_activations" in payload and payload["max_activations"] is not None:
update_data["max_activations"] = int(payload["max_activations"])
if not update_data:
return _error(400, "no_changes")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
promo = await promo_code_dal.update_promo_code(session, promo_id, update_data)
if not promo:
return _error(404, "not_found")
await session.commit()
await session.refresh(promo)
return _ok({"promo": _serialize_promo(promo)})
async def admin_promo_delete_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
promo_id = int(request.match_info["promo_id"])
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
promo = await promo_code_dal.delete_promo_code(session, promo_id)
if not promo:
return _error(404, "not_found")
await session.commit()
return _ok({})
+57
View File
@@ -0,0 +1,57 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
def setup_admin_routes(app: web.Application) -> None:
router = app.router
router.add_get("/api/admin/me", admin_me_route)
router.add_get("/api/admin/stats", admin_stats_route)
router.add_get("/api/admin/users", admin_users_list_route)
router.add_get("/api/admin/users/{user_id:-?\\d+}", admin_user_detail_route)
router.add_get("/api/admin/users/{user_id:-?\\d+}/avatar", admin_user_avatar_route)
router.add_post("/api/admin/users/{user_id:-?\\d+}/ban", admin_user_ban_route)
router.add_post("/api/admin/users/{user_id:-?\\d+}/message", admin_user_message_route)
router.add_post(
"/api/admin/users/{user_id:-?\\d+}/message/preview", admin_user_message_preview_route
)
router.add_post("/api/admin/users/{user_id:-?\\d+}/reset-trial", admin_user_reset_trial_route)
router.add_post("/api/admin/users/{user_id:-?\\d+}/extend", admin_user_extend_route)
router.add_post(
"/api/admin/users/{user_id:-?\\d+}/premium-override",
admin_user_premium_override_route,
)
router.add_post(
"/api/admin/users/{user_id:-?\\d+}/regular-traffic-override",
admin_user_regular_traffic_override_route,
)
router.add_post(
"/api/admin/users/{user_id:-?\\d+}/traffic-grant",
admin_user_traffic_grant_route,
)
router.add_delete("/api/admin/users/{user_id:-?\\d+}", admin_user_delete_route)
router.add_get("/api/admin/payments", admin_payments_list_route)
router.add_get("/api/admin/payments/export.csv", admin_payments_export_route)
router.add_get("/api/admin/promos", admin_promos_list_route)
router.add_post("/api/admin/promos", admin_promo_create_route)
router.add_patch("/api/admin/promos/{promo_id:\\d+}", admin_promo_update_route)
router.add_delete("/api/admin/promos/{promo_id:\\d+}", admin_promo_delete_route)
router.add_get("/api/admin/logs", admin_logs_route)
router.add_post("/api/admin/broadcast", admin_broadcast_route)
router.add_post("/api/admin/sync", admin_sync_route)
router.add_get("/api/admin/ads", admin_ads_list_route)
router.add_post("/api/admin/ads", admin_ad_create_route)
router.add_post("/api/admin/ads/{campaign_id:\\d+}/toggle", admin_ad_toggle_route)
router.add_delete("/api/admin/ads/{campaign_id:\\d+}", admin_ad_delete_route)
router.add_get("/api/admin/settings", admin_settings_get_route)
router.add_patch("/api/admin/settings", admin_settings_patch_route)
router.add_get("/api/admin/tariffs", admin_tariffs_get_route)
router.add_put("/api/admin/tariffs", admin_tariffs_save_route)
router.add_get("/api/admin/panel/internal-squads", admin_panel_internal_squads_route)
+74
View File
@@ -0,0 +1,74 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_settings_get_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
overrides = await app_settings_dal.get_overrides_with_meta(session)
overrides_by_key = {entry["key"]: entry for entry in overrides}
fields = manifest_payload()
sections: Dict[str, Dict[str, Any]] = {}
for field in fields:
key = field["key"]
section_id = field["section"]
if section_id not in sections:
sections[section_id] = {
"id": section_id,
"order": field["section_order"],
"fields": [],
}
override = overrides_by_key.get(key)
value = current_value(settings, key)
is_secret = bool(field.get("secret"))
response_field = {
**field,
"value": "" if is_secret else value,
"overridden": bool(override),
"updated_at": override.get("updated_at") if override else None,
}
if is_secret:
response_field["has_value"] = bool(value)
sections[section_id]["fields"].append(response_field)
ordered_sections = sorted(sections.values(), key=lambda s: s["order"])
return _ok({"sections": ordered_sections})
async def admin_settings_patch_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
payload = await _read_json(request)
updates = payload.get("updates") or {}
deletes = payload.get("deletes") or []
if not isinstance(updates, dict):
return _error(400, "invalid_updates")
if not isinstance(deletes, list):
return _error(400, "invalid_deletes")
result = await update_overrides(
settings,
async_session_factory,
updates=updates,
deletes=deletes,
actor_id=actor_id,
)
if not result.get("ok"):
return web.json_response(
{"ok": False, "error": "validation_failed", "errors": result.get("errors", {})},
status=400,
)
# Bust the public webapp settings cache so users see new values immediately.
cache = request.app.get("webapp_settings_cache")
if isinstance(cache, dict):
cache["ts"] = 0.0
cache["data"] = {}
return _ok({"applied": result.get("applied", 0), "reverted": result.get("reverted", 0)})
+89
View File
@@ -0,0 +1,89 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_me_route(request: web.Request) -> web.Response:
user_id = _require_admin_user_id(request)
settings: Settings = request.app["settings"]
return _ok({}, user_id=user_id, admin_ids=list(settings.ADMIN_IDS or []))
async def admin_stats_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
user_stats = await user_dal.get_enhanced_user_statistics(session)
financial_stats = await payment_dal.get_financial_statistics(session)
sync_status = await panel_sync_dal.get_panel_sync_status(session)
recent_payments = await payment_dal.get_recent_payment_logs_with_user(session, limit=10)
payload = {
"users": user_stats,
"financial": financial_stats,
"panel_sync": {
"status": sync_status.status if sync_status else "never_run",
"last_sync_time": sync_status.last_sync_time.isoformat()
if sync_status and sync_status.last_sync_time
else None,
"details": sync_status.details if sync_status else None,
"users_processed": sync_status.users_processed_from_panel if sync_status else 0,
"subscriptions_synced": sync_status.subscriptions_synced if sync_status else 0,
},
"recent_payments": [_serialize_payment(p) for p in recent_payments],
}
panel_service = request.app.get("panel_service")
if panel_service is not None:
try:
system = await panel_service.get_system_stats()
bandwidth = await panel_service.get_bandwidth_stats()
panel_body: Dict[str, Any] = {
"system": system or {},
"bandwidth": bandwidth or {},
}
try:
nodes = await panel_service.get_nodes_statistics()
panel_body["nodes"] = nodes or {}
except Exception as exc_nodes: # pragma: no cover - optional endpoint
logger.debug("Panel nodes stats unavailable: %s", exc_nodes)
panel_body["nodes"] = {}
try:
today = datetime.now(timezone.utc).date()
start_d = today - timedelta(days=7)
nodes_bw = await panel_service.get_nodes_bandwidth_usage(
start=start_d.isoformat(),
end=today.isoformat(),
top_nodes_limit=64,
)
panel_body["nodes_bandwidth"] = nodes_bw or {}
except Exception as exc_nb: # pragma: no cover - optional endpoint
logger.debug("Panel nodes bandwidth range unavailable: %s", exc_nb)
panel_body["nodes_bandwidth"] = {}
try:
online_map = _panel_nodes_online_by_uuid(panel_body.get("nodes"))
lookups = await panel_service.get_nodes_online_lookups()
for k, v in lookups.get("byUuid", {}).items():
online_map[k] = v
_enrich_bandwidth_nodes_with_online(
panel_body.get("nodes_bandwidth"),
online_map,
lookups.get("byName") or {},
)
except Exception as exc_merge: # pragma: no cover
logger.debug("Panel nodes online merge skipped: %s", exc_merge)
payload["panel"] = panel_body
except Exception as exc:
logger.debug("Panel stats unavailable: %s", exc)
payload["panel"] = {"error": "unavailable"}
queue_manager = get_queue_manager()
if queue_manager:
try:
payload["queue"] = queue_manager.get_queue_stats()
except Exception: # pragma: no cover - defensive
payload["queue"] = None
payload["currency_symbol"] = settings.DEFAULT_CURRENCY_SYMBOL or "RUB"
return _ok(payload)
+23
View File
@@ -0,0 +1,23 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_sync_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
panel_service = request.app.get("panel_service")
if panel_service is None:
return _error(503, "panel_unavailable")
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
i18n = request.app.get("i18n")
from bot.handlers.admin.sync_admin import perform_sync
async with async_session_factory() as session:
result = await perform_sync(
panel_service=panel_service,
session=session,
settings=settings,
i18n_instance=i18n,
)
return _ok({"result": result or {}})
+63
View File
@@ -0,0 +1,63 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_tariffs_get_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
settings: Settings = request.app["settings"]
path = _tariffs_config_path(settings)
try:
config = settings.tariffs_config
except Exception as exc:
logger.warning("Invalid tariffs config requested from admin UI: %s", exc)
return _error(400, "invalid_tariffs_config", str(exc))
if config is None:
return _ok(
{
"exists": path.exists(),
"path": str(path),
"catalog": {
"default_tariff": "",
"topup_packages_default": {"rub": [], "stars": []},
"tariffs": [],
},
}
)
return _ok(
{
"exists": True,
"path": str(path),
"catalog": _tariffs_config_payload(config),
}
)
async def admin_tariffs_save_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
settings: Settings = request.app["settings"]
payload = await _read_json(request)
catalog = payload.get("catalog") if "catalog" in payload else payload
if not isinstance(catalog, dict):
return _error(400, "invalid_payload", "catalog must be an object")
try:
config = TariffsConfig.model_validate(catalog)
except (ValidationError, ValueError) as exc:
return _error(400, "invalid_tariffs_config", str(exc))
path = _tariffs_config_path(settings)
try:
_write_tariffs_config_file(path, config)
except OSError as exc:
logger.exception("Failed to write tariffs config to %s", path)
return _error(500, "write_failed", str(exc))
cache = request.app.get("webapp_settings_cache")
if isinstance(cache, dict):
cache["ts"] = 0.0
cache["data"] = {}
return _ok({"exists": True, "path": str(path), "catalog": _tariffs_config_payload(config)})
+925
View File
@@ -0,0 +1,925 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def admin_users_list_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
page = max(0, int(request.query.get("page", 0) or 0))
page_size = min(100, max(1, int(request.query.get("page_size", 25) or 25)))
query = (request.query.get("q") or "").strip()
filter_value = (request.query.get("filter") or "all").lower()
panel_status = (request.query.get("panel_status") or "all").lower()
premium_traffic = (request.query.get("premium_traffic") or "all").lower()
sort_value = (request.query.get("sort") or "registered_desc").lower()
async with async_session_factory() as session:
users, total = await _filter_and_sort_users(
session,
query=query,
filter_value=filter_value,
panel_status=panel_status,
premium_traffic=premium_traffic,
sort_value=sort_value,
page=page,
page_size=page_size,
)
statuses = await _bulk_user_statuses(session, [u.user_id for u in users])
cached_avatar_ids = await _bulk_user_avatar_keys(session, [u.user_id for u in users])
active_subs = await _bulk_active_subscriptions_for_users(
session, [u.user_id for u in users]
)
serialized = []
for user in users:
payload = _serialize_user(user)
status_payload = statuses.get(user.user_id) or {"status": "bot_only", "end_date": None}
payload["panel_status"] = status_payload.get("status")
if status_payload.get("status") == "expired" and status_payload.get("end_date"):
payload["panel_status_expired_at"] = status_payload["end_date"]
payload["avatar_url"] = (
f"/api/admin/users/{user.user_id}/avatar?v={cached_avatar_ids[user.user_id]}"
if user.user_id in cached_avatar_ids
else None
)
payload["premium_traffic"] = _premium_traffic_list_payload(active_subs.get(user.user_id))
serialized.append(payload)
return _ok(
{
"users": serialized,
"page": page,
"page_size": page_size,
"total": total,
}
)
async def _bulk_user_statuses(
session: AsyncSession, user_ids: List[int]
) -> Dict[int, Dict[str, Optional[str]]]:
"""Return active subscription status for a batch of users.
Returns the panel status (active/expired/limited/disabled) when an active
subscription exists, otherwise ``"bot_only"`` when the user is in the bot
but has no panel subscription.
"""
if not user_ids:
return {}
stmt = (
select(
Subscription.user_id,
Subscription.status_from_panel,
Subscription.is_active,
Subscription.end_date,
)
.where(Subscription.user_id.in_(user_ids))
.order_by(Subscription.is_active.desc(), Subscription.end_date.desc().nullslast())
)
rows = (await session.execute(stmt)).all()
out: Dict[int, Dict[str, Optional[str]]] = {}
for uid, panel_status, is_active, end_date in rows:
if uid in out:
continue
if is_active:
status = (panel_status or "active").lower()
else:
status = (panel_status or "expired").lower()
out[uid] = {
"status": status,
"end_date": end_date.isoformat() if end_date else None,
}
for uid in user_ids:
if uid not in out:
out[uid] = {"status": "bot_only", "end_date": None}
return out
async def _bulk_user_avatar_keys(session: AsyncSession, user_ids: List[int]) -> Dict[int, str]:
"""Return ``{user_id: cache_key}`` for users with a cached Telegram avatar.
The cache key is the row's ``updated_at`` timestamp — used as a
cache-buster query param so the browser refetches when the avatar
changes.
"""
if not user_ids:
return {}
stmt = select(UserTelegramAvatar.user_id, UserTelegramAvatar.updated_at).where(
UserTelegramAvatar.user_id.in_(user_ids)
)
rows = (await session.execute(stmt)).all()
return {int(uid): (updated_at.isoformat() if updated_at else "") for uid, updated_at in rows}
async def admin_user_avatar_route(request: web.Request) -> web.Response:
"""Serve the cached Telegram avatar for any user (admin-only).
Mirrors ``/api/account/avatar`` but takes a ``user_id`` from the URL
and uses admin auth. Only the cached blob from
``user_telegram_avatars`` is served — refreshing from Telegram is the
job of the user-facing endpoint, so the admin list never blocks on a
Telegram round-trip.
"""
_require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
avatar = await session.get(UserTelegramAvatar, target_id)
if not avatar:
raise web.HTTPNotFound(text="avatar_not_cached")
etag = (
f'W/"avatar-{target_id}-{int(avatar.updated_at.timestamp())}"'
if avatar.updated_at
else None
)
if etag and request.headers.get("If-None-Match") == etag:
return web.Response(status=304, headers={"ETag": etag})
response = web.Response(
body=bytes(avatar.image_bytes),
content_type=avatar.content_type or "image/jpeg",
)
response.headers["Cache-Control"] = "private, max-age=3600"
if etag:
response.headers["ETag"] = etag
return response
def _ranked_active_subscriptions_sq(now: datetime):
"""Latest active subscription per user (same ordering as subscription_dal)."""
rn = sa_func.row_number().over(
partition_by=Subscription.user_id,
order_by=(
Subscription.end_date.desc(),
Subscription.subscription_id.desc(),
),
)
inner = (
select(
Subscription.user_id,
Subscription.subscription_id,
Subscription.premium_used_bytes,
Subscription.premium_baseline_bytes,
Subscription.premium_topup_balance_bytes,
Subscription.premium_topup_used_bytes,
Subscription.premium_bonus_bytes,
Subscription.premium_unlimited_override,
rn.label("rn"),
)
.where(
Subscription.is_active.is_(True),
Subscription.end_date > now,
)
.subquery()
)
return select(inner).where(inner.c.rn == 1).subquery(name="ranked_active_sub")
async def _bulk_active_subscriptions_for_users(
session: AsyncSession, user_ids: List[int]
) -> Dict[int, Subscription]:
"""Active subscription row per user (for admin list premium traffic column)."""
if not user_ids:
return {}
now = datetime.now(timezone.utc)
stmt = (
select(Subscription)
.where(
Subscription.user_id.in_(user_ids),
Subscription.is_active.is_(True),
Subscription.end_date > now,
)
.order_by(
Subscription.user_id.asc(),
Subscription.end_date.desc(),
Subscription.subscription_id.desc(),
)
)
rows = (await session.execute(stmt)).scalars().all()
out: Dict[int, Subscription] = {}
for sub in rows:
uid = int(sub.user_id)
if uid not in out:
out[uid] = sub
return out
async def _filter_and_sort_users(
session: AsyncSession,
*,
query: str = "",
filter_value: str,
panel_status: str = "all",
premium_traffic: str = "all",
sort_value: str,
page: int,
page_size: int,
) -> tuple[List[User], int]:
"""Return paginated users with optional search, filter and sort applied."""
now = datetime.now(timezone.utc)
sort_key = (sort_value or "registered_desc").lower()
pt_filter = (premium_traffic or "all").lower()
needs_premium_sq = pt_filter != "all" or sort_key in {
"premium_ratio_asc",
"premium_ratio_desc",
}
stmt = select(User)
count_stmt = select(sa_func.count(User.user_id))
sq = None
ratio_expr = None
plim_expr = None
pu_expr = None
if needs_premium_sq:
sq = _ranked_active_subscriptions_sq(now)
stmt = stmt.outerjoin(sq, User.user_id == sq.c.user_id)
count_stmt = count_stmt.outerjoin(sq, User.user_id == sq.c.user_id)
pb = sa_func.coalesce(sq.c.premium_bonus_bytes, 0)
plim_expr = (
sa_func.coalesce(sq.c.premium_baseline_bytes, 0)
+ sa_func.coalesce(sq.c.premium_topup_balance_bytes, 0)
+ sa_func.coalesce(sq.c.premium_topup_used_bytes, 0)
+ pb
)
pu_expr = sa_func.coalesce(sq.c.premium_used_bytes, 0)
ratio_expr = case(
(sq.c.user_id.is_(None), None),
(sq.c.premium_unlimited_override.is_(True), None),
(plim_expr <= 0, None),
else_=cast(pu_expr, Float) / cast(plim_expr, Float),
)
search_cond = _user_search_condition(query)
if search_cond is not None:
stmt = stmt.where(search_cond)
count_stmt = count_stmt.where(search_cond)
f = (filter_value or "all").lower()
if f == "banned":
cond = User.is_banned.is_(True)
elif f == "active":
cond = User.is_banned.is_(False)
elif f == "tg_linked":
cond = User.telegram_id.is_not(None)
elif f == "no_tg":
cond = User.telegram_id.is_(None)
elif f == "email_linked":
cond = User.email.is_not(None)
elif f == "no_email":
cond = User.email.is_(None)
elif f == "panel_linked":
cond = User.panel_user_uuid.is_not(None)
else:
cond = None
if cond is not None:
stmt = stmt.where(cond)
count_stmt = count_stmt.where(cond)
panel_cond = _user_panel_status_condition(panel_status)
if panel_cond is not None:
stmt = stmt.where(panel_cond)
count_stmt = count_stmt.where(panel_cond)
if needs_premium_sq and sq is not None and plim_expr is not None and pu_expr is not None:
if pt_filter == "none":
premium_cond = or_(
sq.c.user_id.is_(None),
and_(
sq.c.premium_unlimited_override.is_(False),
plim_expr <= 0,
),
)
stmt = stmt.where(premium_cond)
count_stmt = count_stmt.where(premium_cond)
elif pt_filter == "unlimited":
premium_cond = and_(
sq.c.user_id.isnot(None),
sq.c.premium_unlimited_override.is_(True),
)
stmt = stmt.where(premium_cond)
count_stmt = count_stmt.where(premium_cond)
elif pt_filter == "good":
premium_cond = and_(
sq.c.user_id.isnot(None),
sq.c.premium_unlimited_override.is_(False),
plim_expr > 0,
(100 * pu_expr) < (85 * plim_expr),
)
stmt = stmt.where(premium_cond)
count_stmt = count_stmt.where(premium_cond)
elif pt_filter == "warn":
premium_cond = and_(
sq.c.user_id.isnot(None),
sq.c.premium_unlimited_override.is_(False),
plim_expr > 0,
(100 * pu_expr) >= (85 * plim_expr),
pu_expr < plim_expr,
)
stmt = stmt.where(premium_cond)
count_stmt = count_stmt.where(premium_cond)
elif pt_filter == "critical":
premium_cond = and_(
sq.c.user_id.isnot(None),
sq.c.premium_unlimited_override.is_(False),
plim_expr > 0,
pu_expr >= plim_expr,
)
stmt = stmt.where(premium_cond)
count_stmt = count_stmt.where(premium_cond)
sort_map = {
"registered_desc": User.registration_date.desc().nullslast(),
"registered_asc": User.registration_date.asc().nullslast(),
"name_asc": (
sa_func.coalesce(User.first_name, User.username, User.email).asc(),
User.user_id.asc(),
),
"name_desc": (
sa_func.coalesce(User.first_name, User.username, User.email).desc(),
User.user_id.desc(),
),
"id_asc": User.user_id.asc(),
"id_desc": User.user_id.desc(),
}
if needs_premium_sq and ratio_expr is not None and sort_key == "premium_ratio_asc":
stmt = stmt.order_by(ratio_expr.asc().nullslast(), User.user_id.asc())
elif needs_premium_sq and ratio_expr is not None and sort_key == "premium_ratio_desc":
stmt = stmt.order_by(ratio_expr.desc().nullslast(), User.user_id.desc())
else:
order = sort_map.get(sort_key, sort_map["registered_desc"])
if isinstance(order, tuple):
stmt = stmt.order_by(*order)
else:
stmt = stmt.order_by(order)
stmt = stmt.offset(max(page, 0) * max(page_size, 1)).limit(max(page_size, 1))
users = (await session.execute(stmt)).scalars().all()
total = (await session.execute(count_stmt)).scalar_one()
return users, int(total)
def _user_panel_status_condition(panel_status: str):
status = (panel_status or "all").lower()
if status not in {"active", "expired", "limited"}:
return None
normalized_status = sa_func.lower(sa_func.coalesce(Subscription.status_from_panel, ""))
blank_status = or_(
Subscription.status_from_panel.is_(None), Subscription.status_from_panel == ""
)
if status == "active":
status_cond = or_(
normalized_status == "active", blank_status & Subscription.is_active.is_(True)
)
elif status == "expired":
status_cond = or_(
normalized_status == "expired", blank_status & Subscription.is_active.is_(False)
)
else:
status_cond = normalized_status == "limited"
return (
select(Subscription.subscription_id)
.where(Subscription.user_id == User.user_id, status_cond)
.exists()
)
def _user_search_condition(query: str):
raw = (query or "").strip().lstrip("@")
if not raw:
return None
like = f"%{raw}%"
conditions = [
User.username.ilike(like),
User.first_name.ilike(like),
User.last_name.ilike(like),
User.email.ilike(like),
]
if raw.isdigit():
numeric = int(raw)
conditions.extend([User.user_id == numeric, User.telegram_id == numeric])
return or_(*conditions)
async def admin_user_detail_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
async_session_factory: sessionmaker = request.app["async_session_factory"]
settings: Settings = request.app["settings"]
async with async_session_factory() as session:
user = await user_dal.get_user_by_id(session, target_id)
if not user:
return _error(404, "not_found", "User not found")
active_sub = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
latest_subs_stmt = (
select(Subscription)
.where(Subscription.user_id == target_id)
.order_by(Subscription.start_date.desc().nullslast())
.limit(20)
)
latest_subs = (await session.execute(latest_subs_stmt)).scalars().all()
total_paid = await payment_dal.get_user_total_paid(session, target_id)
recent_payments_stmt = (
select(Payment)
.where(Payment.user_id == target_id)
.order_by(Payment.created_at.desc())
.limit(20)
)
recent_payments = (await session.execute(recent_payments_stmt)).scalars().all()
log_count = await message_log_dal.count_user_message_logs(session, target_id)
avatar_keys = await _bulk_user_avatar_keys(session, [target_id])
# Referral links — both the bot deep-link and the webapp deep-link.
referral_code: Optional[str] = None
try:
referral_code = await user_dal.ensure_referral_code(session, user)
await session.commit()
except Exception as exc_ref: # pragma: no cover — defensive
logger.warning("Failed to ensure referral code for user %s: %s", target_id, exc_ref)
await session.rollback()
referral_service: Optional[ReferralService] = request.app.get("referral_service")
bot_username = request.app.get("bot_username") or ""
referral_bot_link: Optional[str] = None
if referral_service and bot_username and referral_code:
try:
async with async_session_factory() as session:
referral_bot_link = await referral_service.generate_referral_link(
session, bot_username, target_id
)
except Exception as exc_link: # pragma: no cover
logger.warning("Failed to build bot referral link for %s: %s", target_id, exc_link)
referral_webapp_link = _build_admin_webapp_referral_link(
getattr(settings, "SUBSCRIPTION_MINI_APP_URL", None),
referral_code,
)
# Subscription page URL — the raw panel `subscriptionUrl` that the user
# imports into their VPN client. May be missing if the user has never
# been provisioned on the panel.
subscription_url: Optional[str] = None
panel_uuid = getattr(user, "panel_user_uuid", None)
if panel_uuid:
subscription_service = request.app.get("subscription_service")
panel_service = getattr(subscription_service, "panel_service", None)
if panel_service is not None:
try:
panel_data = await panel_service.get_user_by_uuid(panel_uuid)
if panel_data:
subscription_url = panel_data.get("subscriptionUrl") or None
except Exception as exc_panel: # pragma: no cover
logger.warning(
"Failed to fetch subscriptionUrl for user %s (uuid=%s): %s",
target_id,
panel_uuid,
exc_panel,
)
serialized_user = _serialize_user(user)
serialized_user["avatar_url"] = (
f"/api/admin/users/{target_id}/avatar?v={avatar_keys[target_id]}"
if target_id in avatar_keys
else None
)
return _ok(
{
"user": serialized_user,
"active_subscription": _serialize_subscription(active_sub) if active_sub else None,
"subscriptions": [_serialize_subscription(s) for s in (latest_subs or [])],
"total_paid": float(total_paid),
"recent_payments": [_serialize_payment(p) for p in recent_payments],
"log_count": int(log_count or 0),
"subscription_url": subscription_url,
"referral": {
"code": referral_code,
"bot_link": referral_bot_link,
"webapp_link": referral_webapp_link,
},
}
)
async def admin_user_ban_route(request: web.Request) -> web.Response:
_require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
desired = bool(payload.get("banned"))
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
user = await user_dal.get_user_by_id(session, target_id)
if not user:
return _error(404, "not_found")
user.is_banned = bool(desired)
await session.commit()
await session.refresh(user)
return _ok({"user": _serialize_user(user)})
async def admin_user_message_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
text = str(payload.get("text") or "").strip()
if not text:
return _error(400, "empty_text")
queue_manager = get_queue_manager()
if not queue_manager:
return _error(503, "queue_unavailable")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
target_user = await user_dal.get_user_by_id(session, target_id)
if not target_user or not target_user.telegram_id:
return _error(404, "no_telegram_account")
try:
await send_message_via_queue(
queue_manager,
int(target_user.telegram_id),
MessageContent(content_type="text", text=text),
parse_mode="HTML",
disable_web_page_preview=True,
)
except Exception as exc:
logger.warning("Admin direct message failed: %s", exc)
return _error(502, "send_failed", str(exc))
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_direct_message_webapp",
"content": text[:4000],
"is_admin_event": True,
"target_user_id": target_id,
},
)
return _ok({})
async def admin_user_message_preview_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
admin_telegram_id = request.get("admin_telegram_id")
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
text = str(payload.get("text") or "").strip()
if not text:
return _error(400, "empty_text")
if not admin_telegram_id:
return _error(403, "admin_telegram_unavailable")
queue_manager = get_queue_manager()
if not queue_manager:
return _error(503, "queue_unavailable")
try:
await send_message_via_queue(
queue_manager,
int(admin_telegram_id),
MessageContent(content_type="text", text=text),
parse_mode="HTML",
disable_web_page_preview=True,
)
except Exception as exc:
logger.warning("Admin direct message preview failed: %s", exc)
return _error(502, "preview_failed", str(exc))
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_direct_message_preview_webapp",
"content": text[:4000],
"is_admin_event": True,
"target_user_id": target_id,
},
)
return _ok({})
async def admin_user_delete_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
ok = await user_dal.delete_user_and_relations(session, target_id)
if not ok:
await session.rollback()
return _error(404, "not_found")
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_delete_user_webapp",
"content": f"Deleted user_id={target_id}",
"is_admin_event": True,
},
)
await session.commit()
return _ok({})
async def admin_user_reset_trial_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
panel_service = request.app.get("panel_service")
subscription_service = request.app.get("subscription_service")
if panel_service is None or subscription_service is None:
return _error(503, "service_unavailable")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
user = await user_dal.get_user_by_id(session, target_id)
if not user:
return _error(404, "not_found")
active = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
if active:
await session.delete(active)
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_reset_trial_webapp",
"content": f"Reset trial for user_id={target_id}",
"is_admin_event": True,
"target_user_id": target_id,
},
)
await session.commit()
return _ok({})
async def admin_user_premium_override_route(request: web.Request) -> web.Response:
"""Premium-squad traffic overrides only (unlimited toggle + bonus GB)."""
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
subscription_service = request.app.get("subscription_service")
unlimited = bool(payload.get("unlimited"))
bonus_bytes_raw = payload.get("bonus_bytes")
bonus_gb_raw = payload.get("bonus_gb")
if bonus_bytes_raw is None and bonus_gb_raw is None:
bonus_bytes = 0
elif bonus_bytes_raw is not None:
try:
bonus_bytes = int(bonus_bytes_raw)
except (TypeError, ValueError):
return _error(400, "invalid_bonus", "bonus_bytes must be an integer")
else:
try:
bonus_bytes = int(round(float(bonus_gb_raw) * (1024**3)))
except (TypeError, ValueError):
return _error(400, "invalid_bonus", "bonus_gb must be a number")
if bonus_bytes < 0:
return _error(400, "invalid_bonus", "bonus must be non-negative")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
active = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
if not active:
return _error(404, "no_active_subscription")
active.premium_unlimited_override = bool(unlimited)
active.premium_bonus_bytes = int(bonus_bytes)
if active.premium_unlimited_override:
active.premium_is_limited = False
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_premium_override_webapp",
"content": (f"unlimited={bool(unlimited)} bonus_bytes={int(bonus_bytes)}"),
"is_admin_event": True,
"target_user_id": target_id,
},
)
await session.commit()
await session.refresh(active)
if subscription_service is not None:
await subscription_service.sync_premium_squad_access_to_panel(session, target_id)
await session.commit()
await session.refresh(active)
return _ok({"subscription": _serialize_subscription(active)})
async def admin_user_regular_traffic_override_route(request: web.Request) -> web.Response:
"""Main (regular) traffic: unlimited-style ceiling + admin bonus GB."""
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
unlimited = bool(payload.get("unlimited"))
regular_bonus_bytes_raw = payload.get("regular_bonus_bytes")
regular_bonus_gb_raw = payload.get("regular_bonus_gb")
if regular_bonus_bytes_raw is None and regular_bonus_gb_raw is None:
regular_bonus_bytes = 0
elif regular_bonus_bytes_raw is not None:
try:
regular_bonus_bytes = int(regular_bonus_bytes_raw)
except (TypeError, ValueError):
return _error(400, "invalid_regular_bonus", "regular_bonus_bytes must be an integer")
else:
try:
regular_bonus_bytes = int(round(float(regular_bonus_gb_raw) * (1024**3)))
except (TypeError, ValueError):
return _error(400, "invalid_regular_bonus", "regular_bonus_gb must be a number")
if regular_bonus_bytes < 0:
return _error(400, "invalid_regular_bonus", "regular bonus must be non-negative")
subscription_service = request.app.get("subscription_service")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
active = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
if not active:
return _error(404, "no_active_subscription")
active.regular_unlimited_override = bool(unlimited)
active.regular_bonus_bytes = int(regular_bonus_bytes)
if subscription_service is not None:
await subscription_service.sync_main_traffic_limit_to_panel(session, target_id)
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_regular_traffic_override_webapp",
"content": (
f"unlimited={bool(unlimited)} regular_bonus_bytes={int(regular_bonus_bytes)}"
),
"is_admin_event": True,
"target_user_id": target_id,
},
)
await session.commit()
await session.refresh(active)
return _ok({"subscription": _serialize_subscription(active)})
async def admin_user_traffic_grant_route(request: web.Request) -> web.Response:
"""Credit regular or premium traffic to a user without a payment.
Body: ``{"kind": "regular" | "premium", "gb": float}`` (alternatively
``"bytes": int``). Mirrors the same effect as a user-purchased top-up:
the chosen balance grows, panel limit/squads are refreshed, and an entry
is added to ``traffic_topups`` with ``kind="admin_topup"`` or
``kind="admin_premium_topup"`` and ``payment_id=NULL``.
"""
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
kind = str(payload.get("kind") or "regular").strip().lower()
if kind not in {"regular", "premium"}:
return _error(400, "invalid_kind", "kind must be 'regular' or 'premium'")
bytes_raw = payload.get("bytes")
gb_raw = payload.get("gb")
if bytes_raw is None and gb_raw is None:
return _error(400, "missing_amount", "either 'gb' or 'bytes' is required")
try:
if bytes_raw is not None:
grant_bytes = int(bytes_raw)
gb_value = grant_bytes / (1024**3)
else:
gb_value = float(gb_raw)
grant_bytes = int(round(gb_value * (1024**3)))
except (TypeError, ValueError):
return _error(400, "invalid_amount", "amount must be a positive number")
if gb_value <= 0 or grant_bytes <= 0:
return _error(400, "invalid_amount", "amount must be positive")
subscription_service = request.app.get("subscription_service")
if subscription_service is None:
return _error(503, "subscription_service_unavailable")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
active = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
if not active:
return _error(404, "no_active_subscription")
if kind == "regular":
result = await subscription_service.admin_grant_topup(session, target_id, gb_value)
else:
result = await subscription_service.admin_grant_premium_topup(
session, target_id, gb_value
)
if not result:
await session.rollback()
return _error(
422,
"grant_failed",
"Unable to credit traffic (missing tariff/squads or panel error)",
)
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_traffic_grant_webapp",
"content": f"kind={kind} bytes={grant_bytes}",
"is_admin_event": True,
"target_user_id": target_id,
},
)
await session.commit()
refreshed = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
return _ok(
{
"subscription": _serialize_subscription(refreshed) if refreshed else None,
"grant": {
"kind": kind,
"granted_bytes": grant_bytes,
"granted_gb": gb_value,
},
}
)
async def admin_user_extend_route(request: web.Request) -> web.Response:
actor_id = _require_admin_user_id(request)
target_id = int(request.match_info["user_id"])
payload = await _read_json(request)
try:
days = int(payload.get("days") or 0)
except (TypeError, ValueError):
return _error(400, "invalid_days")
if days <= 0:
return _error(400, "invalid_days")
subscription_service = request.app.get("subscription_service")
if subscription_service is None:
return _error(503, "subscription_service_unavailable")
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
new_end = await subscription_service.extend_active_subscription_days(
session,
target_id,
days,
"admin_extend_subscription_webapp",
)
if not new_end:
await session.rollback()
return _error(500, "extend_failed")
await message_log_dal.create_message_log(
session,
{
"user_id": actor_id,
"event_type": "admin_extend_subscription_webapp",
"content": f"+{days}d -> {new_end.isoformat()}",
"is_admin_event": True,
"target_user_id": target_id,
},
)
await session.commit()
refreshed = await subscription_dal.get_active_subscription_by_user_id(session, target_id)
return _ok(
{
"subscription": _serialize_subscription(refreshed) if refreshed else None,
}
)
+29
View File
@@ -0,0 +1,29 @@
from __future__ import annotations
from typing import Optional
from aiohttp import web
from bot.app.web.webapp_auth import verify_webapp_session_token
from config.settings import Settings
WEBAPP_SESSION_COOKIE_NAME = "rw_webapp_session"
def extract_authenticated_user_id(request: web.Request) -> Optional[int]:
settings: Settings = request.app["settings"]
auth_header = request.headers.get("Authorization", "")
if auth_header.startswith("Bearer "):
user_id = verify_webapp_session_token(
settings,
auth_header.removeprefix("Bearer ").strip(),
)
if user_id:
return user_id
session_cookie = request.cookies.get(WEBAPP_SESSION_COOKIE_NAME)
if session_cookie:
return verify_webapp_session_token(settings, session_cookie)
return None
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -72,7 +72,7 @@ async def build_and_start_web_app(
secret_token=settings.WEBHOOK_SECRET_TOKEN,
).register(app, path=telegram_webhook_path)
logging.info(
f"Telegram webhook route configured at: [POST] {telegram_webhook_path} (relative to base URL)"
f"Telegram webhook route configured at: [POST] {telegram_webhook_path} (relative to base URL)" # noqa: E501
)
from bot.handlers.user.payment import yookassa_webhook_route
+1
View File
@@ -0,0 +1 @@
"""Domain modules for the subscription Mini App backend."""
+100
View File
@@ -0,0 +1,100 @@
# ruff: noqa: F401,F403,F405,I001
import asyncio
import base64
import hashlib
import hmac
import io
import ipaddress
import json
import logging
import os
import re
import secrets
import socket
import subprocess
import time
from collections import deque
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
from aiogram import Bot, Dispatcher
from aiogram.types import LabeledPrice
from aiohttp import ClientSession, ClientTimeout, web
from pydantic import BaseModel, ConfigDict, EmailStr, ValidationError, constr, field_validator
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import sessionmaker
from bot.app.web.admin_api import (
admin_auth_middleware,
setup_admin_routes,
)
from bot.app.web.webapp_auth import (
create_signed_telegram_oauth_state,
create_telegram_oauth_nonce,
create_webapp_session_token,
validate_telegram_login_widget_data,
validate_telegram_oauth_id_token,
validate_telegram_webapp_init_data,
verify_signed_telegram_oauth_state,
verify_telegram_oauth_nonce,
verify_webapp_session_token,
)
from bot.services.crypto_pay_service import CryptoPayService
from bot.services.email_auth_service import EmailAuthService, normalize_email
from bot.services.email_templates import render_account_merged
from bot.services.freekassa_service import FreeKassaService
from bot.services.platega_service import PlategaService
from bot.services.promo_code_service import PromoCodeService
from bot.services.referral_service import ReferralService
from bot.services.severpay_service import SeverPayService
from bot.services.subscription_service import SubscriptionService
from bot.services.yookassa_service import YooKassaService
from bot.utils.config_link import prepare_config_links
from bot.utils.request_security import parse_ip_entries, request_client_ip
from bot.utils.text_sanitizer import sanitize_display_name, sanitize_username
from config.settings import Settings
from db.dal import payment_dal, subscription_dal, user_dal
from db.dal.user_dal import UserMergeConflictError
from db.models import Payment, User, UserTelegramAvatar
logger = logging.getLogger(__name__)
TEMPLATE_PATH = Path(__file__).resolve().parents[1] / "templates" / "subscription_webapp.html"
ASSET_DIR = TEMPLATE_PATH.parent
WEBAPP_LOGO_PROXY_PATH = "/webapp-logo"
WEBAPP_LOGO_CACHE_DIR = Path(__file__).resolve().parents[4] / "data" / "webapp-logo"
WEBAPP_EMOJI_CACHE_DIR = Path(__file__).resolve().parents[4] / "data" / "webapp-emoji"
WEBAPP_CONFIG_PLACEHOLDER = "<!-- WEBAPP_CONFIG_SCRIPT -->"
WEBAPP_I18N_PLACEHOLDER = "<!-- WEBAPP_I18N_SCRIPT -->"
WEBAPP_JS_PLACEHOLDER = "<!-- WEBAPP_JS_SCRIPT -->"
APP_REPOSITORY_URL = "https://github.com/3252a8/remnawave-minishop"
DEV_MOCK_START_MARKER = "<!-- WEBAPP_DEV_MOCK_START -->"
DEV_MOCK_END_MARKER = "<!-- WEBAPP_DEV_MOCK_END -->"
WEBAPP_RATE_LIMIT_WINDOW_SECONDS = 60
WEBAPP_RATE_LIMIT_MAX_REQUESTS = 30
WEBAPP_LOGO_MAX_BYTES = 2 * 1024 * 1024
WEBAPP_EMOJI_MAX_BYTES = 4 * 1024 * 1024
WEBAPP_TELEGRAM_AVATAR_MAX_BYTES = 128 * 1024
WEBAPP_TELEGRAM_AVATAR_REFRESH_SECONDS = 24 * 60 * 60
WEBAPP_TELEGRAM_AVATAR_FETCH_TIMEOUT_SECONDS = 4
WEBAPP_SESSION_COOKIE_NAME = "rw_webapp_session"
WEBAPP_CSRF_COOKIE_NAME = "rw_webapp_csrf"
WEBAPP_TELEGRAM_OAUTH_STATE_COOKIE_NAME = "rw_tg_oauth_state"
WEBAPP_CSRF_HEADER_NAME = "X-CSRF-Token"
WEBAPP_STATE_CHANGING_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
_APP_VERSION_CACHE: Optional[str] = None
WEBAPP_CSRF_EXEMPT_PATHS = {
"/api/auth/telegram/nonce",
"/api/auth/token",
"/api/auth/email/request",
"/api/auth/email/verify",
"/api/auth/email/magic",
"/api/auth/logout",
}
_SHARED_HTTP_SESSION: Optional[ClientSession] = None
_SHARED_HTTP_SESSION_LOCK = asyncio.Lock()
__all__ = [name for name in globals() if not name.startswith("__")]
+534
View File
@@ -0,0 +1,534 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def account_email_request_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
settings: Settings = request.app["settings"]
payload = await _read_json(request)
email_payload, validation_error = _validate_model_payload(WebAppEmailPayload, payload)
if validation_error:
return validation_error
email = email_payload.email
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
return _json_error(403, "access_denied", "Access denied")
if db_user.email == email and db_user.email_verified_at:
return web.json_response({"ok": True, "already_linked": True})
lang = _normalize_language(db_user.language_code or settings.DEFAULT_LANGUAGE)
return await _request_email_code(
request,
email=email,
purpose="link_email",
language_code=lang,
target_user_id=user_id,
)
async def account_email_verify_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
rate_limit_response = await _enforce_webapp_rate_limit(
request,
user_id=user_id,
action="account_email_verify",
)
if rate_limit_response:
return rate_limit_response
payload = await _read_json(request)
email_payload, validation_error = _validate_model_payload(WebAppEmailCodePayload, payload)
if validation_error:
return validation_error
email = email_payload.email
code = str(email_payload.code or "")
email_service: EmailAuthService = request.app["email_auth_service"]
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
merge_notice: Optional[Dict[str, Any]] = None
source_panel_uuid: Optional[str] = None
final_user_id = user_id
final_email = email
final_telegram_id: Optional[int] = None
final_username: Optional[str] = None
final_first_name: Optional[str] = None
final_panel_uuid: Optional[str] = None
should_notify_email_linked = False
async with async_session_factory() as session:
try:
verify_result = await email_service.verify_code(
session,
email=email,
purpose="link_email",
code=code,
target_user_id=user_id,
)
if not verify_result.ok:
await session.commit()
status = 429 if verify_result.error == "rate_limited" else 400
return web.json_response(
{
"ok": False,
"error": verify_result.error or "invalid_code",
"retry_after": verify_result.retry_after,
"message": "Invalid code",
},
status=status,
)
current_user = await user_dal.get_user_by_id(session, user_id)
if not current_user or current_user.is_banned:
await session.rollback()
return _json_error(403, "access_denied", "Access denied")
should_notify_email_linked = (
bool(_telegram_id_for_user(current_user)) and not current_user.email
)
existing_email_user = await user_dal.get_user_by_email(session, email)
if existing_email_user and existing_email_user.user_id != current_user.user_id:
source_panel_uuid = existing_email_user.panel_user_uuid
current_user = await user_dal.merge_users(
session,
source_user_id=existing_email_user.user_id,
target_user_id=current_user.user_id,
)
merge_notice = await _build_account_merge_notice(
session,
merged_user=current_user,
source_user_id=existing_email_user.user_id,
source_panel_uuid=source_panel_uuid,
settings=settings,
)
current_user.email = email
current_user.email_verified_at = datetime.now(timezone.utc)
await _sync_panel_identity_for_user(request, current_user)
await session.commit()
final_user_id = int(current_user.user_id)
final_telegram_id = _telegram_id_for_user(current_user)
final_username = current_user.username
final_first_name = current_user.first_name
final_panel_uuid = current_user.panel_user_uuid
if merge_notice:
merge_end_date_raw = merge_notice.get("final_end_date")
merge_end_date = (
datetime.fromisoformat(merge_end_date_raw) if merge_end_date_raw else None
)
await _sync_panel_identity_for_user(
request,
current_user,
expire_at=merge_end_date,
)
# Best-effort cleanup of the removed panel account after the DB merge.
if source_panel_uuid and final_panel_uuid and source_panel_uuid != final_panel_uuid:
subscription_service: SubscriptionService = request.app.get(
"subscription_service"
)
if subscription_service and subscription_service.panel_service:
try:
await subscription_service.panel_service.delete_user_from_panel(
source_panel_uuid,
log_response=False,
)
except Exception as exc:
logger.warning(
"Failed to delete merged source panel user %s: %s",
source_panel_uuid,
exc,
)
email_service: EmailAuthService = request.app.get("email_auth_service")
if email_service and final_email:
email_content = render_account_merged(
settings,
language_code=merge_notice.get("language") or settings.DEFAULT_LANGUAGE,
primary_user_id=merge_notice.get("primary_user_id"),
removed_user_id=merge_notice.get("removed_user_id"),
final_end_date_text=str(
merge_notice.get("final_end_date_text")
or merge_notice.get("final_end_date")
or ""
),
)
try:
await email_service.send_rendered_email(
email=final_email,
content=email_content,
)
except Exception as exc:
logger.warning(
"Failed to send account merge email to %s: %s",
final_email,
exc,
)
except UserMergeConflictError as exc:
await session.rollback()
return _json_error(409, "account_merge_conflict", str(exc))
except Exception:
await session.rollback()
logger.exception("Email account link failed")
return _json_error(500, "link_failed", "Link failed")
if should_notify_email_linked:
try:
from bot.services.notification_service import NotificationService
bot: Bot = request.app["bot"]
notification_service = NotificationService(
bot,
settings,
request.app.get("i18n"),
)
await notification_service.notify_account_email_linked(
user_id=int(final_user_id),
email=final_email,
telegram_id=final_telegram_id,
username=final_username,
first_name=final_first_name,
)
except Exception:
logger.exception("Failed to send account email linked notification")
token = create_webapp_session_token(settings, int(final_user_id))
response_payload: Dict[str, Any] = {"ok": True}
if merge_notice:
response_payload["account_merge"] = merge_notice
response_payload["user_id"] = final_user_id
return _build_webapp_auth_response(settings, response_payload, token=token)
async def account_telegram_link_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
settings: Settings = request.app["settings"]
payload = await _read_json(request)
telegram_user = await _validate_telegram_auth_payload(request, payload)
if not telegram_user:
return _json_error(401, "invalid_auth", "Invalid Telegram auth data")
async_session_factory: sessionmaker = request.app["async_session_factory"]
merge_notice: Optional[Dict[str, Any]] = None
source_panel_uuid: Optional[str] = None
final_user_id = user_id
final_telegram_id: Optional[int] = None
final_email: Optional[str] = None
final_username: Optional[str] = None
final_first_name: Optional[str] = None
final_panel_uuid: Optional[str] = None
should_notify_telegram_linked = False
async with async_session_factory() as session:
try:
current_user_before_link = await user_dal.get_user_by_id(session, user_id)
if not current_user_before_link or current_user_before_link.is_banned:
await session.rollback()
return _json_error(403, "access_denied", "Access denied")
should_notify_telegram_linked = bool(
current_user_before_link.email
) and not _telegram_id_for_user(current_user_before_link)
source_panel_uuid = current_user_before_link.panel_user_uuid
db_user = await _link_telegram_to_user(
request,
session,
current_user_id=user_id,
telegram_user=telegram_user,
settings=settings,
)
if db_user.is_banned:
await session.rollback()
return _json_error(403, "banned", "Access denied")
final_user_id = int(db_user.user_id)
final_telegram_id = _telegram_id_for_user(db_user)
final_email = db_user.email
final_username = db_user.username
final_first_name = db_user.first_name
final_panel_uuid = db_user.panel_user_uuid
if final_user_id != user_id:
merge_notice = await _build_account_merge_notice(
session,
merged_user=db_user,
source_user_id=user_id,
source_panel_uuid=source_panel_uuid,
settings=settings,
)
await session.commit()
if merge_notice:
merge_end_date_raw = merge_notice.get("final_end_date")
merge_end_date = (
datetime.fromisoformat(merge_end_date_raw) if merge_end_date_raw else None
)
await _sync_panel_identity_for_user(
request,
db_user,
expire_at=merge_end_date,
)
# Best-effort cleanup of the removed panel account after the DB merge.
if source_panel_uuid and final_panel_uuid and source_panel_uuid != final_panel_uuid:
subscription_service: SubscriptionService = request.app.get(
"subscription_service"
)
if subscription_service and subscription_service.panel_service:
try:
await subscription_service.panel_service.delete_user_from_panel(
source_panel_uuid,
log_response=False,
)
except Exception as exc:
logger.warning(
"Failed to delete merged source panel user %s: %s",
source_panel_uuid,
exc,
)
email_service: EmailAuthService = request.app.get("email_auth_service")
if email_service and final_email:
email_content = render_account_merged(
settings,
language_code=merge_notice.get("language") or settings.DEFAULT_LANGUAGE,
primary_user_id=merge_notice.get("primary_user_id"),
removed_user_id=merge_notice.get("removed_user_id"),
final_end_date_text=str(
merge_notice.get("final_end_date_text")
or merge_notice.get("final_end_date")
or ""
),
)
try:
await email_service.send_rendered_email(
email=final_email,
content=email_content,
)
except Exception as exc:
logger.warning(
"Failed to send account merge email to %s: %s",
final_email,
exc,
)
except UserMergeConflictError as exc:
await session.rollback()
return _json_error(409, "account_merge_conflict", str(exc))
except Exception:
await session.rollback()
logger.exception("Telegram account link failed")
return _json_error(500, "link_failed", "Link failed")
if should_notify_telegram_linked and final_telegram_id:
try:
from bot.services.notification_service import NotificationService
bot: Bot = request.app["bot"]
notification_service = NotificationService(
bot,
settings,
request.app.get("i18n"),
)
await notification_service.notify_account_telegram_linked(
user_id=int(final_user_id),
email=final_email,
telegram_id=int(final_telegram_id),
username=final_username,
first_name=final_first_name,
)
except Exception:
logger.exception("Failed to send account Telegram linked notification")
token = create_webapp_session_token(settings, int(final_user_id))
response_payload: Dict[str, Any] = {
"ok": True,
"user_id": int(final_user_id),
"telegram_id": final_telegram_id,
}
if merge_notice:
response_payload["account_merge"] = merge_notice
return _build_webapp_auth_response(settings, response_payload, token=token)
async def me_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
data = await _build_user_payload(request, user_id)
return web.json_response({"ok": True, **data})
async def account_avatar_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
await session.rollback()
return _json_error(403, "access_denied", "Access denied")
avatar = await _ensure_cached_telegram_avatar(request, session, db_user)
await session.commit()
if not avatar:
raise web.HTTPNotFound(text="avatar_not_cached")
etag = _telegram_avatar_etag(avatar)
if etag and request.headers.get("If-None-Match") == etag:
return web.Response(status=304, headers={"ETag": etag})
response = web.Response(
body=bytes(avatar.image_bytes),
content_type=avatar.content_type or "image/jpeg",
)
response.headers["Cache-Control"] = "private, max-age=3600"
if etag:
response.headers["ETag"] = etag
return response
async def account_language_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
payload = await _read_json(request)
language_payload, validation_error = _validate_model_payload(WebAppLanguagePayload, payload)
if validation_error:
return validation_error
language = _normalize_language(str(language_payload.language or ""))
async_session_factory: sessionmaker = request.app["async_session_factory"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
await session.rollback()
return _json_error(403, "access_denied", "Access denied")
if _normalize_language(db_user.language_code or "") != language:
db_user.language_code = language
await session.flush()
await session.commit()
return web.json_response({"ok": True, "language": language})
def _format_webapp_datetime(value: Optional[datetime]) -> Optional[str]:
if not value:
return None
normalized = value if value.tzinfo else value.replace(tzinfo=timezone.utc)
return normalized.strftime("%d.%m.%Y %H:%M")
def _telegram_photo_url_value(telegram_user: Dict[str, Any]) -> Optional[str]:
raw_value = telegram_user.get("photo_url")
if not raw_value:
return None
value = str(raw_value).strip()
return value or None
def _telegram_avatar_is_stale(avatar: Optional[UserTelegramAvatar]) -> bool:
if not avatar or not avatar.updated_at:
return True
updated_at = avatar.updated_at
if updated_at.tzinfo is None:
updated_at = updated_at.replace(tzinfo=timezone.utc)
return (
datetime.now(timezone.utc) - updated_at
).total_seconds() >= WEBAPP_TELEGRAM_AVATAR_REFRESH_SECONDS
def _telegram_avatar_etag(avatar: UserTelegramAvatar) -> str:
digest = hashlib.sha256(bytes(avatar.image_bytes)).hexdigest()[:16]
return f'"tg-avatar-{int(avatar.user_id)}-{digest}"'
def _telegram_avatar_url(avatar: Optional[UserTelegramAvatar]) -> str:
if not avatar:
return ""
updated_at = avatar.updated_at
if updated_at and updated_at.tzinfo is None:
updated_at = updated_at.replace(tzinfo=timezone.utc)
version = (
int(updated_at.timestamp())
if updated_at
else hashlib.sha256(bytes(avatar.image_bytes)).hexdigest()[:8]
)
return f"/api/account/avatar?v={version}"
def _select_compact_telegram_photo_size(sizes: List[Any]) -> Optional[Any]:
if not sizes:
return None
suitable = [size for size in sizes if int(getattr(size, "width", 0) or 0) >= 160]
candidates = suitable or sizes
return min(
candidates,
key=lambda size: (
int(getattr(size, "file_size", 0) or 0)
or int(getattr(size, "width", 0) or 0) * int(getattr(size, "height", 0) or 0),
int(getattr(size, "width", 0) or 0),
),
)
def _telegram_file_content_type(file_path: Optional[str]) -> str:
path = str(file_path or "").lower()
if path.endswith(".png"):
return "image/png"
if path.endswith(".webp"):
return "image/webp"
return "image/jpeg"
async def _fetch_compact_telegram_avatar(
bot: Bot, telegram_id: int
) -> Optional[Tuple[bytes, str, Optional[str]]]:
photos = await bot.get_user_profile_photos(user_id=telegram_id, limit=1)
if not photos or not photos.photos:
return None
photo_size = _select_compact_telegram_photo_size(list(photos.photos[0] or []))
if not photo_size:
return None
file_info = await bot.get_file(photo_size.file_id)
destination = io.BytesIO()
await bot.download_file(file_info.file_path, destination=destination)
body = destination.getvalue()
if not body or len(body) > WEBAPP_TELEGRAM_AVATAR_MAX_BYTES:
return None
return (
body,
_telegram_file_content_type(file_info.file_path),
getattr(photo_size, "file_unique_id", None),
)
async def _ensure_cached_telegram_avatar(
request: web.Request,
session: AsyncSession,
user: User,
) -> Optional[UserTelegramAvatar]:
avatar = await user_dal.get_user_telegram_avatar(session, int(user.user_id))
telegram_id = _telegram_id_for_user(user)
if not telegram_id:
return avatar
if avatar and not _telegram_avatar_is_stale(avatar):
return avatar
bot: Bot = request.app["bot"]
try:
fetched = await asyncio.wait_for(
_fetch_compact_telegram_avatar(bot, int(telegram_id)),
timeout=WEBAPP_TELEGRAM_AVATAR_FETCH_TIMEOUT_SECONDS,
)
except Exception as exc:
logger.info("Failed to refresh Telegram avatar for user %s: %s", user.user_id, exc)
return avatar
if not fetched:
return avatar
body, content_type, file_unique_id = fetched
return await user_dal.upsert_user_telegram_avatar(
session,
user_id=int(user.user_id),
file_unique_id=file_unique_id,
content_type=content_type,
image_bytes=body,
)
+60
View File
@@ -0,0 +1,60 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
def create_subscription_webapp_application(
dp: Dispatcher,
bot: Bot,
settings: Settings,
async_session_factory: sessionmaker,
) -> web.Application:
app = web.Application(
middlewares=[
_security_headers_middleware,
_csrf_protection_middleware,
admin_auth_middleware,
]
)
app["bot"] = bot
app["dp"] = dp
app["settings"] = settings
app["async_session_factory"] = async_session_factory
app["i18n"] = dp.get("i18n_instance")
app["email_auth_service"] = EmailAuthService(settings)
app["webapp_logo_cache"] = None
app["webapp_logo_cache_lock"] = asyncio.Lock()
app["webapp_settings_cache"] = {"ts": 0.0, "data": {}}
app["webapp_rate_limit_buckets"] = {}
app["webapp_rate_limit_lock"] = asyncio.Lock()
async def _startup(app_obj: web.Application) -> None:
await _ensure_shared_http_session()
await _warm_webapp_logo_cache(app_obj)
await _warm_webapp_animated_emoji_cache(app_obj)
async def _shutdown(app_obj: web.Application) -> None:
await _close_shared_http_session()
app.on_startup.append(_startup)
app.on_shutdown.append(_shutdown)
for key in (
"subscription_service",
"yookassa_service",
"freekassa_service",
"cryptopay_service",
"platega_service",
"severpay_service",
"promo_code_service",
"referral_service",
"panel_service",
):
if hasattr(dp, "workflow_data") and key in dp.workflow_data: # type: ignore[attr-defined]
app[key] = dp.workflow_data[key] # type: ignore[index]
# type: ignore[attr-defined]
if hasattr(dp, "workflow_data") and "bot_username" in dp.workflow_data:
app["bot_username"] = dp.workflow_data["bot_username"] # type: ignore[index]
setup_subscription_webapp_routes(app)
return app
+734
View File
@@ -0,0 +1,734 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def health_route(request: web.Request) -> web.Response:
return web.json_response({"ok": True})
async def css_asset_route(request: web.Request) -> web.Response:
return await _serve_template_asset(request, "subscription_webapp.css", "text/css")
def _resolve_webapp_logo_url(settings: Settings) -> str:
raw_logo_url = (settings.WEBAPP_LOGO_URL or "").strip()
if not raw_logo_url:
return ""
parsed_logo_url = urlsplit(raw_logo_url)
if parsed_logo_url.scheme == "https":
cache_key = hashlib.sha256(raw_logo_url.encode("utf-8")).hexdigest()[:12]
return f"{WEBAPP_LOGO_PROXY_PATH}?v={cache_key}"
if parsed_logo_url.scheme in {"http", "data"}:
return raw_logo_url
if raw_logo_url.startswith("/"):
return raw_logo_url
return ""
def _webapp_logo_cache_key(logo_url: str) -> str:
return hashlib.sha256(logo_url.encode("utf-8")).hexdigest()
def _webapp_logo_disk_paths(logo_url: str) -> Tuple[Path, Path]:
cache_key = _webapp_logo_cache_key(logo_url)
return WEBAPP_LOGO_CACHE_DIR / f"{cache_key}.bin", WEBAPP_LOGO_CACHE_DIR / f"{cache_key}.json"
def _is_proxyable_webapp_logo_url(logo_url: str) -> bool:
parsed_logo_url = urlsplit(logo_url)
return parsed_logo_url.scheme == "https" and bool(parsed_logo_url.hostname)
def _emoji_to_codepoints(value: str) -> str:
return "_".join(f"{ord(char):x}" for char in str(value or "").strip())
def _webapp_emoji_disk_path(codepoints: str, ext: str) -> Path:
return WEBAPP_EMOJI_CACHE_DIR / f"{codepoints}.512.{ext}"
def _webapp_animated_emoji_source_url(codepoints: str, ext: str) -> str:
return f"https://fonts.gstatic.com/s/e/notoemoji/latest/{codepoints}/512.{ext}"
def _webapp_animated_emoji_asset_path(emoji: str, ext: str = "gif") -> str:
codepoints = _emoji_to_codepoints(emoji)
if not codepoints or ext not in {"gif", "webp"}:
return ""
return f"/webapp-emoji/{codepoints}/512.{ext}"
async def webapp_logo_route(request: web.Request) -> web.Response:
settings: Settings = request.app["settings"]
raw_logo_url = (settings.WEBAPP_LOGO_URL or "").strip()
if not raw_logo_url:
raise web.HTTPNotFound(text="webapp_logo_not_configured")
if not _is_proxyable_webapp_logo_url(raw_logo_url):
raise web.HTTPNotFound(text="webapp_logo_not_proxied")
parsed_logo_url = urlsplit(raw_logo_url)
if not await _hostname_resolves_to_public_address(parsed_logo_url.hostname):
raise web.HTTPNotFound(text="webapp_logo_not_proxied")
source_logo_url = raw_logo_url
logo_cache: Optional[Tuple[str, bytes, str]] = request.app.get("webapp_logo_cache")
if logo_cache is None or logo_cache[0] != source_logo_url:
cache_lock: asyncio.Lock = request.app["webapp_logo_cache_lock"]
async with cache_lock:
logo_cache = request.app.get("webapp_logo_cache")
if logo_cache is None or logo_cache[0] != source_logo_url:
fetched_logo = await _load_or_fetch_webapp_logo(source_logo_url)
logo_cache = (
(source_logo_url, fetched_logo[0], fetched_logo[1]) if fetched_logo else None
)
request.app["webapp_logo_cache"] = logo_cache
if not logo_cache:
raise web.HTTPNotFound(text="webapp_logo_unavailable")
_, body, content_type = logo_cache
response = web.Response(body=body, content_type=content_type)
response.headers["Cache-Control"] = "public, max-age=31536000, immutable"
return response
async def webapp_animated_emoji_route(request: web.Request) -> web.Response:
codepoints = str(request.match_info.get("codepoints") or "").strip().lower()
ext = str(request.match_info.get("ext") or "").strip().lower()
if not re.fullmatch(r"[0-9a-f]+(?:_[0-9a-f]+)*", codepoints) or ext not in {"gif", "webp"}:
raise web.HTTPNotFound(text="webapp_emoji_not_found")
emoji_cache_key = f"{codepoints}:{ext}"
emoji_caches: Dict[str, Tuple[bytes, str]] = request.app.setdefault("webapp_emoji_cache", {})
emoji_cache = emoji_caches.get(emoji_cache_key)
if emoji_cache is None:
cache_lock: asyncio.Lock = request.app.setdefault("webapp_emoji_cache_lock", asyncio.Lock())
async with cache_lock:
emoji_cache = emoji_caches.get(emoji_cache_key)
if emoji_cache is None:
emoji_cache = await _load_or_fetch_webapp_animated_emoji(codepoints, ext)
if emoji_cache:
emoji_caches[emoji_cache_key] = emoji_cache
if not emoji_cache:
raise web.HTTPNotFound(text="webapp_emoji_unavailable")
body, content_type = emoji_cache
response = web.Response(body=body, content_type=content_type)
response.headers["Cache-Control"] = "public, max-age=31536000, immutable"
return response
async def _warm_webapp_logo_cache(app: web.Application) -> None:
settings: Settings = app["settings"]
raw_logo_url = (settings.WEBAPP_LOGO_URL or "").strip()
if not raw_logo_url or not _is_proxyable_webapp_logo_url(raw_logo_url):
return
parsed_logo_url = urlsplit(raw_logo_url)
if not parsed_logo_url.hostname or not await _hostname_resolves_to_public_address(
parsed_logo_url.hostname
):
return
cache_lock: asyncio.Lock = app["webapp_logo_cache_lock"]
async with cache_lock:
logo_cache: Optional[Tuple[str, bytes, str]] = app.get("webapp_logo_cache")
if logo_cache and logo_cache[0] == raw_logo_url:
return
loaded_logo = await _load_or_fetch_webapp_logo(raw_logo_url)
app["webapp_logo_cache"] = (
(raw_logo_url, loaded_logo[0], loaded_logo[1]) if loaded_logo else None
)
async def _warm_webapp_animated_emoji_cache(app: web.Application) -> None:
settings: Settings = app["settings"]
if str(settings.WEBAPP_LOGO_EMOJI_FONT or "").strip() != "noto-color-animated":
return
codepoints = _emoji_to_codepoints(settings.WEBAPP_LOGO_EMOJI)
if not codepoints:
return
app.setdefault("webapp_emoji_cache", {})
app.setdefault("webapp_emoji_cache_lock", asyncio.Lock())
emoji_caches: Dict[str, Tuple[bytes, str]] = app["webapp_emoji_cache"]
for ext in ("gif", "webp"):
emoji_cache_key = f"{codepoints}:{ext}"
if emoji_cache_key in emoji_caches:
continue
loaded_emoji = await _load_or_fetch_webapp_animated_emoji(codepoints, ext)
if loaded_emoji:
emoji_caches[emoji_cache_key] = loaded_emoji
if ext == "gif":
return
async def _load_or_fetch_webapp_animated_emoji(
codepoints: str, ext: str
) -> Optional[Tuple[bytes, str]]:
disk_emoji = await asyncio.to_thread(_read_webapp_animated_emoji_from_disk, codepoints, ext)
if disk_emoji:
return disk_emoji
fetched_emoji = await _fetch_webapp_animated_emoji(codepoints, ext)
if fetched_emoji:
await asyncio.to_thread(
_write_webapp_animated_emoji_to_disk, codepoints, ext, fetched_emoji
)
return fetched_emoji
def _read_webapp_animated_emoji_from_disk(codepoints: str, ext: str) -> Optional[Tuple[bytes, str]]:
path = _webapp_emoji_disk_path(codepoints, ext)
try:
body = path.read_bytes()
except OSError:
return None
if not body or len(body) > WEBAPP_EMOJI_MAX_BYTES:
return None
return body, "image/gif" if ext == "gif" else "image/webp"
def _write_webapp_animated_emoji_to_disk(
codepoints: str, ext: str, emoji: Tuple[bytes, str]
) -> None:
body, _content_type = emoji
if not body or len(body) > WEBAPP_EMOJI_MAX_BYTES:
return
path = _webapp_emoji_disk_path(codepoints, ext)
try:
WEBAPP_EMOJI_CACHE_DIR.mkdir(parents=True, exist_ok=True)
path.write_bytes(body)
except OSError as exc:
logger.warning("Failed to write WEBAPP animated emoji cache: %s", exc)
async def _fetch_webapp_animated_emoji(codepoints: str, ext: str) -> Optional[Tuple[bytes, str]]:
try:
session = await _get_shared_http_session()
timeout = ClientTimeout(total=4)
source_url = _webapp_animated_emoji_source_url(codepoints, ext)
async with session.get(
source_url,
allow_redirects=False,
headers={"Accept": "image/gif,image/webp,image/*,*/*;q=0.8"},
timeout=timeout,
) as response:
if response.status != 200:
return None
content_type = (
(response.headers.get("Content-Type") or "").split(";", 1)[0].strip().lower()
)
expected_content_type = "image/gif" if ext == "gif" else "image/webp"
if content_type and content_type != expected_content_type:
return None
body = bytearray()
async for chunk in response.content.iter_chunked(64 * 1024):
body.extend(chunk)
if len(body) > WEBAPP_EMOJI_MAX_BYTES:
logger.warning("WEBAPP animated emoji exceeded the 4 MiB limit.")
return None
if not body:
return None
return bytes(body), expected_content_type
except Exception as exc:
logger.warning("Failed to fetch WEBAPP animated emoji: %s", exc)
return None
async def _load_or_fetch_webapp_logo(logo_url: str) -> Optional[Tuple[bytes, str]]:
disk_logo = await asyncio.to_thread(_read_webapp_logo_from_disk, logo_url)
if disk_logo:
return disk_logo
fetched_logo = await _fetch_webapp_logo(logo_url)
if fetched_logo:
await asyncio.to_thread(_write_webapp_logo_to_disk, logo_url, fetched_logo)
return fetched_logo
def _read_webapp_logo_from_disk(logo_url: str) -> Optional[Tuple[bytes, str]]:
body_path, meta_path = _webapp_logo_disk_paths(logo_url)
try:
metadata = json.loads(meta_path.read_text(encoding="utf-8"))
if metadata.get("source_url") != logo_url:
return None
content_type = str(metadata.get("content_type") or "").strip().lower()
if not content_type.startswith("image/"):
return None
body = body_path.read_bytes()
except (OSError, json.JSONDecodeError):
return None
if not body or len(body) > WEBAPP_LOGO_MAX_BYTES:
return None
return body, content_type
def _write_webapp_logo_to_disk(logo_url: str, logo: Tuple[bytes, str]) -> None:
body, content_type = logo
if not body or len(body) > WEBAPP_LOGO_MAX_BYTES:
return
body_path, meta_path = _webapp_logo_disk_paths(logo_url)
try:
WEBAPP_LOGO_CACHE_DIR.mkdir(parents=True, exist_ok=True)
body_path.write_bytes(body)
meta_path.write_text(
json.dumps(
{
"source_url": logo_url,
"content_type": content_type,
"cached_at": datetime.now(timezone.utc).isoformat(),
"bytes": len(body),
},
ensure_ascii=False,
separators=(",", ":"),
),
encoding="utf-8",
)
except OSError as exc:
logger.warning("Failed to write WEBAPP_LOGO_URL cache: %s", exc)
async def _fetch_webapp_logo(logo_url: str) -> Optional[Tuple[bytes, str]]:
"""Fetch and cache the configured logo on the server side."""
try:
session = await _get_shared_http_session()
timeout = ClientTimeout(total=3)
async with session.get(
logo_url,
allow_redirects=False,
headers={"Accept": "image/avif,image/webp,image/svg+xml,image/png,image/*,*/*;q=0.8"},
timeout=timeout,
) as response:
if response.status != 200:
logger.warning(
"WEBAPP_LOGO_URL returned HTTP %s; keeping the logo hidden.",
response.status,
)
return None
content_type = (
(response.headers.get("Content-Type") or "").split(";", 1)[0].strip().lower()
)
if content_type and not content_type.startswith("image/"):
logger.warning(
"WEBAPP_LOGO_URL returned non-image content type %s; keeping the logo hidden.",
content_type,
)
return None
body = bytearray()
async for chunk in response.content.iter_chunked(64 * 1024):
body.extend(chunk)
if len(body) > WEBAPP_LOGO_MAX_BYTES:
logger.warning("WEBAPP_LOGO_URL exceeded the 2 MiB limit.")
return None
if not body:
logger.warning("WEBAPP_LOGO_URL returned an empty response body.")
return None
return bytes(body), content_type or "image/png"
except Exception as exc:
logger.warning("Failed to fetch WEBAPP_LOGO_URL: %s", exc)
return None
async def _get_shared_http_session() -> ClientSession:
global _SHARED_HTTP_SESSION
async with _SHARED_HTTP_SESSION_LOCK:
if _SHARED_HTTP_SESSION is None or _SHARED_HTTP_SESSION.closed:
_SHARED_HTTP_SESSION = ClientSession(
timeout=ClientTimeout(total=30),
headers={
"User-Agent": "Mozilla/5.0",
"Accept": "*/*",
},
)
return _SHARED_HTTP_SESSION
async def _ensure_shared_http_session() -> None:
await _get_shared_http_session()
async def _close_shared_http_session() -> None:
global _SHARED_HTTP_SESSION
async with _SHARED_HTTP_SESSION_LOCK:
if _SHARED_HTTP_SESSION and not _SHARED_HTTP_SESSION.closed:
await _SHARED_HTTP_SESSION.close()
_SHARED_HTTP_SESSION = None
async def _hostname_resolves_to_public_address(hostname: str) -> bool:
if not hostname:
return False
try:
ip_obj = ipaddress.ip_address(hostname)
return not (
ip_obj.is_private
or ip_obj.is_loopback
or ip_obj.is_link_local
or ip_obj.is_unspecified
or ip_obj.is_reserved
)
except ValueError:
pass
loop = asyncio.get_running_loop()
try:
resolved = await loop.getaddrinfo(hostname, None, type=socket.SOCK_STREAM)
except Exception:
return False
found_public_ip = False
for entry in resolved:
sockaddr = entry[4]
if not sockaddr:
continue
candidate = sockaddr[0]
try:
ip_obj = ipaddress.ip_address(candidate)
except ValueError:
continue
if (
ip_obj.is_private
or ip_obj.is_loopback
or ip_obj.is_link_local
or ip_obj.is_unspecified
or ip_obj.is_reserved
):
return False
found_public_ip = True
return found_public_ip
@web.middleware
async def _security_headers_middleware(request: web.Request, handler):
request["csp_nonce"] = secrets.token_urlsafe(16)
try:
response = await handler(request)
except web.HTTPException as exc:
response = exc
nonce = request.get("csp_nonce", "")
response.headers.setdefault(
"Content-Security-Policy",
(
"default-src 'self'; "
f"script-src 'self' 'nonce-{nonce}' https://telegram.org; "
"frame-src https://oauth.telegram.org; "
"frame-ancestors https://web.telegram.org https://t.me; "
"style-src 'self' 'unsafe-inline' https://fonts.googleapis.com https://cdn.jsdelivr.net; " # noqa: E501
"font-src 'self' https://fonts.gstatic.com https://cdn.jsdelivr.net data:; "
"img-src 'self' data: https:; "
"connect-src 'self' https://oauth.telegram.org; "
"object-src 'none'; "
"base-uri 'self'; "
"form-action 'self'"
),
)
response.headers.setdefault("Referrer-Policy", "no-referrer")
response.headers.setdefault("X-Content-Type-Options", "nosniff")
response.headers.setdefault(
"Permissions-Policy",
(
"accelerometer=(), autoplay=(), camera=(), display-capture=(), "
"encrypted-media=(), geolocation=(), gyroscope=(), magnetometer=(), "
"microphone=(), midi=(), payment=(), usb=()"
),
)
return response
@web.middleware
async def _csrf_protection_middleware(request: web.Request, handler):
settings: Settings = request.app["settings"]
header = request.headers.get("Authorization", "")
prefix = "Bearer "
if header.startswith(prefix):
if verify_webapp_session_token(settings, header[len(prefix) :].strip()):
return await handler(request)
if (
request.method in WEBAPP_STATE_CHANGING_METHODS
and request.path not in WEBAPP_CSRF_EXEMPT_PATHS
and request.cookies.get(WEBAPP_SESSION_COOKIE_NAME)
):
csrf_cookie = request.cookies.get(WEBAPP_CSRF_COOKIE_NAME, "")
csrf_header = request.headers.get(WEBAPP_CSRF_HEADER_NAME, "")
if not csrf_cookie or not csrf_header or not hmac.compare_digest(csrf_header, csrf_cookie):
return _json_error(403, "csrf_failed", "Invalid CSRF token")
return await handler(request)
def _get_cached_webapp_settings(request: web.Request) -> Dict[str, Any]:
settings: Settings = request.app["settings"]
cache = request.app["webapp_settings_cache"]
now = time.monotonic()
if now - float(cache.get("ts", 0.0)) >= 60 or not cache.get("data"):
cache["data"] = {
"logo_url": _resolve_webapp_logo_url(settings),
"subscription_options": settings.subscription_options,
"stars_subscription_options": settings.stars_subscription_options,
"traffic_packages": settings.traffic_packages,
"stars_traffic_packages": settings.stars_traffic_packages,
"support_url": settings.SUPPORT_LINK or "",
"terms_url": settings.TERMS_OF_SERVICE_URL or "",
"privacy_policy_url": settings.PRIVACY_POLICY_URL or "",
"user_agreement_url": settings.USER_AGREEMENT_URL or "",
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"email_auth_enabled": settings.email_auth_configured,
"language": _normalize_language(settings.DEFAULT_LANGUAGE),
}
cache["ts"] = now
return cache["data"]
def _run_git_command(*args: str) -> str:
repo_root = Path(__file__).resolve().parents[3]
try:
result = subprocess.run(
["git", *args],
cwd=repo_root,
check=True,
capture_output=True,
text=True,
timeout=1.5,
)
except (OSError, subprocess.SubprocessError):
return ""
return result.stdout.strip()
def _resolve_app_version() -> str:
global _APP_VERSION_CACHE
if _APP_VERSION_CACHE:
return _APP_VERSION_CACHE
env_version = os.getenv("REMNAWAVE_MINISHOP_VERSION", "").strip()
if env_version:
_APP_VERSION_CACHE = env_version
return env_version
build_version_path = Path(__file__).resolve().parents[3] / ".build-version"
try:
build_version = build_version_path.read_text(encoding="utf-8").strip()
except OSError:
build_version = ""
if build_version:
_APP_VERSION_CACHE = build_version
return build_version
tag = _run_git_command("describe", "--tags", "--abbrev=0")
sha = _run_git_command("rev-parse", "--short", "HEAD")
dirty = bool(_run_git_command("status", "--porcelain"))
if tag and sha:
commits_since_tag = _run_git_command("rev-list", f"{tag}..HEAD", "--count")
if commits_since_tag and commits_since_tag != "0":
version = f"{tag}+{commits_since_tag}.g{sha}"
else:
version = tag
elif sha:
version = f"dev+g{sha}"
else:
version = "dev+unknown"
if dirty:
version = f"{version}-dirty"
_APP_VERSION_CACHE = version
return version
async def _enforce_webapp_rate_limit(
request: web.Request,
*,
user_id: int,
action: str,
) -> Optional[web.Response]:
settings: Settings = request.app["settings"]
ip_address = (
request_client_ip(request, trusted_proxies=settings.trusted_proxies)
or request.remote
or "unknown"
)
key = f"{action}:{ip_address}:{int(user_id)}"
buckets: Dict[str, deque[float]] = request.app["webapp_rate_limit_buckets"]
lock: asyncio.Lock = request.app["webapp_rate_limit_lock"]
now = time.monotonic()
async with lock:
bucket = buckets.setdefault(key, deque())
while bucket and now - bucket[0] >= WEBAPP_RATE_LIMIT_WINDOW_SECONDS:
bucket.popleft()
if not bucket:
buckets.pop(key, None)
bucket = buckets.setdefault(key, deque())
if len(bucket) >= WEBAPP_RATE_LIMIT_MAX_REQUESTS:
retry_after = (
max(
1,
int(WEBAPP_RATE_LIMIT_WINDOW_SECONDS - (now - bucket[0])),
)
if bucket
else WEBAPP_RATE_LIMIT_WINDOW_SECONDS
)
return web.json_response(
{
"ok": False,
"error": "rate_limited",
"retry_after": retry_after,
},
status=429,
headers={"Retry-After": str(retry_after)},
)
bucket.append(now)
return None
async def js_asset_route(request: web.Request) -> web.Response:
asset_hash = request.match_info.get("asset_hash")
filename = (
f"subscription_webapp.min.{asset_hash}.js" if asset_hash else "subscription_webapp.js"
)
response = await _serve_template_asset(
request,
filename,
"application/javascript",
strip_dev_mock=not asset_hash,
)
response.headers["Cache-Control"] = (
"public, max-age=31536000, immutable" if asset_hash else "no-cache"
)
return response
async def index_route(request: web.Request) -> web.Response:
settings: Settings = request.app["settings"]
if not settings.WEBAPP_ENABLED:
raise web.HTTPNotFound(text="webapp_disabled")
html = TEMPLATE_PATH.read_text(encoding="utf-8")
cached = _get_cached_webapp_settings(request)
config = {
"title": settings.WEBAPP_TITLE,
"primaryColor": settings.WEBAPP_PRIMARY_COLOR,
"logoUrl": cached["logo_url"],
"logoEmoji": settings.WEBAPP_LOGO_EMOJI,
"logoEmojiFont": settings.WEBAPP_LOGO_EMOJI_FONT,
"apiBase": "/api",
"telegramLoginBotUsername": request.app.get("bot_username") or "",
"telegramLoginBotId": _resolve_telegram_bot_id(settings.BOT_TOKEN) or 0,
"telegramOAuthClientId": _resolve_telegram_oauth_client_id(settings) or 0,
"telegramOAuthRequestAccess": _resolve_telegram_oauth_request_access(settings),
"supportUrl": cached["support_url"],
"termsUrl": cached["terms_url"],
"privacyPolicyUrl": cached["privacy_policy_url"],
"userAgreementUrl": cached["user_agreement_url"],
"currency": cached["currency"],
"language": cached["language"],
"emailAuthEnabled": cached["email_auth_enabled"],
"appVersion": _resolve_app_version(),
"appRepositoryUrl": APP_REPOSITORY_URL,
}
html = _strip_marked_block(html, DEV_MOCK_START_MARKER, DEV_MOCK_END_MARKER)
i18n_instance: Optional[object] = request.app.get("i18n")
i18n_payload = getattr(i18n_instance, "locales_data", {}) if i18n_instance else {}
nonce = request.get("csp_nonce", "")
html = html.replace(
WEBAPP_CONFIG_PLACEHOLDER,
(
f'<script id="webapp-config" type="application/json" nonce="{nonce}">'
+ json.dumps(config, ensure_ascii=False, separators=(",", ":"))
+ "</script>"
),
)
html = html.replace(
WEBAPP_I18N_PLACEHOLDER,
(
f'<script id="i18n" type="application/json" nonce="{nonce}">'
+ json.dumps(i18n_payload, ensure_ascii=False, separators=(",", ":"))
+ "</script>"
),
)
html = html.replace(
WEBAPP_JS_PLACEHOLDER,
f'<script src="/{_resolve_webapp_js_asset_name()}" defer></script>',
)
brand_asset_url = cached["logo_url"]
if not brand_asset_url and settings.WEBAPP_LOGO_EMOJI_FONT == "noto-color-animated":
brand_asset_url = _webapp_animated_emoji_asset_path(settings.WEBAPP_LOGO_EMOJI)
if brand_asset_url:
html = html.replace(
'<link rel="preload" id="logo-preload" href="" as="image" fetchpriority="high" crossorigin="anonymous">', # noqa: E501
f'<link rel="preload" href="{brand_asset_url}" as="image" fetchpriority="high" crossorigin="anonymous">', # noqa: E501
)
else:
html = html.replace(
'<link rel="preload" id="logo-preload" href="" as="image" fetchpriority="high" crossorigin="anonymous">', # noqa: E501
"",
)
return web.Response(text=html, content_type="text/html", charset="utf-8")
async def _serve_template_asset(
request: web.Request,
filename: str,
content_type: str,
*,
strip_dev_mock: bool = False,
) -> web.Response:
settings: Settings = request.app["settings"]
if not settings.WEBAPP_ENABLED:
raise web.HTTPNotFound(text="webapp_disabled")
path = ASSET_DIR / filename
text = path.read_text(encoding="utf-8")
if strip_dev_mock:
text = _strip_marked_block(
text,
"/* WEBAPP_DEV_MOCK_START */",
"/* WEBAPP_DEV_MOCK_END */",
)
return web.Response(text=text, content_type=content_type, charset="utf-8")
def _resolve_webapp_js_asset_name() -> str:
minified_assets = []
for path in ASSET_DIR.glob("subscription_webapp.min.*.js"):
try:
minified_assets.append((path.stat().st_mtime, path.name))
except OSError:
continue
if minified_assets:
minified_assets.sort(reverse=True)
return minified_assets[0][1]
return "subscription_webapp.js"
def _strip_marked_block(html: str, start_marker: str, end_marker: str) -> str:
start = html.find(start_marker)
if start == -1:
return html
end = html.find(end_marker, start)
if end == -1:
return html[:start]
return html[:start] + html[end + len(end_marker) :]
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+156
View File
@@ -0,0 +1,156 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def _read_json(request: web.Request) -> Dict[str, Any]:
try:
data = await request.json()
return data if isinstance(data, dict) else {}
except Exception:
return {}
def _json_error(status: int, code: str, message: str) -> web.Response:
return web.json_response(
{"ok": False, "error": code, "message": message},
status=status,
)
def _validation_error_response(exc: ValidationError) -> web.Response:
for error in exc.errors():
loc = error.get("loc") or ()
field = str(loc[0]) if loc else ""
error_type = str(error.get("type") or "")
message = str(error.get("msg") or "")
message_lower = message.lower()
if field == "email":
if (
"too_long" in message_lower
or "too long" in message_lower
or error_type == "string_too_long"
):
return _json_error(400, "email_too_long", "Email is too long")
return _json_error(400, "invalid_email", "Invalid email")
if field in {"description", "comment", "note"} and error_type == "string_too_long":
return _json_error(400, f"{field}_too_long", f"{field.capitalize()} is too long")
if error_type == "string_too_long":
return _json_error(400, "text_too_long", "Text is too long")
return _json_error(400, "invalid_request", "Invalid request")
def _validate_model_payload(
model_cls: type[BaseModel],
payload: Dict[str, Any],
) -> tuple[Optional[BaseModel], Optional[web.Response]]:
try:
return model_cls.model_validate(payload), None
except ValidationError as exc:
return None, _validation_error_response(exc)
def _normalize_language(lang: Optional[str]) -> str:
value = (lang or "ru").split("-")[0].lower()
return value if value in {"ru", "en"} else "ru"
def _format_remaining(seconds: int, lang: str) -> str:
if seconds <= 0:
if lang == "en":
return "Subscription inactive"
return "Подписка не активна"
days, rem = divmod(seconds, 86400)
hours, rem = divmod(rem, 3600)
minutes = rem // 60
if lang == "en":
if days > 0:
return f"{days} d. {hours} h."
if hours > 0:
return f"{hours} h. {minutes} min."
return f"{max(1, minutes)} min."
if days > 0:
return f"{days} д. {hours} ч."
if hours > 0:
return f"{hours} ч. {minutes} мин."
return f"{max(1, minutes)} мин."
def _coerce_int_or_none(value: Optional[Any]) -> Optional[int]:
if value is None:
return None
try:
return int(value)
except (TypeError, ValueError):
return None
def _format_bytes(value: Optional[Any], *, zero_as_unlimited: bool = False) -> str:
if value is None:
return "N/A"
try:
size = float(value)
except (TypeError, ValueError):
return str(value)
if size <= 0 and zero_as_unlimited:
return ""
if size <= 0:
size = 0
units = ["B", "KB", "MB", "GB", "TB"]
index = 0
while size >= 1024 and index < len(units) - 1:
size /= 1024
index += 1
return f"{size:.2f} {units[index]}"
def _format_months_title(months: int, lang: str) -> str:
if lang == "en":
if months == 1:
return "1 month"
return f"{months} months"
if months == 1:
return "1 месяц"
if 2 <= months <= 4:
return f"{months} месяца"
return f"{months} месяцев"
def _format_number_for_payload(value: Any) -> str:
numeric = float(value or 0)
return str(int(numeric)) if numeric.is_integer() else f"{numeric:g}"
def _format_traffic_title(traffic_gb: float, lang: str) -> str:
return f"{_format_number_for_payload(traffic_gb)} GB"
def _traffic_payment_description(traffic_gb: float, lang: str) -> str:
if lang == "en":
return f"Traffic package {_format_traffic_title(traffic_gb, lang)}"
return f"Пакет трафика {_format_traffic_title(traffic_gb, lang)}"
def _hwid_devices_payment_description(device_count: int, lang: str) -> str:
if lang == "en":
return f"HWID device package +{device_count}"
return f"Докупка устройств HWID +{device_count}"
def _resolve_numeric_option_key(options: Dict[Any, Any], target: float) -> Optional[Any]:
for key in options:
try:
if abs(float(key) - float(target)) < 0.000001:
return key
except (TypeError, ValueError):
continue
return None
def _payment_description(months: int, lang: str) -> str:
if lang == "en":
return f"Subscription for {_format_months_title(months, lang)}"
return f"Подписка на {_format_months_title(months, lang)}"
+169
View File
@@ -0,0 +1,169 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def devices_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
settings: Settings = request.app["settings"]
if not settings.MY_DEVICES_SECTION_ENABLED:
return _json_error(404, "devices_disabled", "Devices section is disabled")
async_session_factory: sessionmaker = request.app["async_session_factory"]
subscription_service: SubscriptionService = request.app["subscription_service"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
return _json_error(403, "access_denied", "Access denied")
active = await subscription_service.get_active_subscription_details(session, user_id)
panel_user_uuid = active.get("user_id") if active else None
if not panel_user_uuid:
return _json_error(400, "subscription_not_active", "Subscription is not active")
panel_service = getattr(subscription_service, "panel_service", None)
if not panel_service:
return _json_error(503, "panel_unavailable", "Panel service unavailable")
try:
devices_response = await panel_service.get_user_devices(panel_user_uuid)
except Exception:
logger.exception("Failed to load WebApp devices for user %s", user_id)
return _json_error(502, "devices_load_failed", "Failed to load devices")
devices = _normalize_devices_response(devices_response)
max_devices = _coerce_int_or_none(active.get("max_devices")) if active else None
return web.json_response(
{
"ok": True,
"enabled": True,
"current_devices": len(devices),
"max_devices": max_devices,
"max_devices_label": _format_devices_limit(max_devices),
"devices": [
_serialize_device(device, index) for index, device in enumerate(devices, start=1)
],
}
)
async def disconnect_device_route(request: web.Request) -> web.Response:
user_id = _require_user_id(request)
rate_limit_response = await _enforce_webapp_rate_limit(
request,
user_id=user_id,
action="devices_disconnect",
)
if rate_limit_response:
return rate_limit_response
settings: Settings = request.app["settings"]
if not settings.MY_DEVICES_SECTION_ENABLED:
return _json_error(404, "devices_disabled", "Devices section is disabled")
payload = await _read_json(request)
disconnect_payload, validation_error = _validate_model_payload(
WebAppDeviceDisconnectPayload, payload
)
if validation_error:
return validation_error
token = str(disconnect_payload.token or "").strip()
async_session_factory: sessionmaker = request.app["async_session_factory"]
subscription_service: SubscriptionService = request.app["subscription_service"]
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
return _json_error(403, "access_denied", "Access denied")
active = await subscription_service.get_active_subscription_details(session, user_id)
panel_user_uuid = active.get("user_id") if active else None
if not panel_user_uuid:
return _json_error(400, "subscription_not_active", "Subscription is not active")
panel_service = getattr(subscription_service, "panel_service", None)
if not panel_service:
return _json_error(503, "panel_unavailable", "Panel service unavailable")
try:
devices_response = await panel_service.get_user_devices(panel_user_uuid)
except Exception:
logger.exception("Failed to load WebApp devices before disconnect for user %s", user_id)
return _json_error(502, "devices_load_failed", "Failed to load devices")
target_hwid = None
for device in _normalize_devices_response(devices_response):
hwid = str(device.get("hwid") or "").strip()
if hwid and hmac.compare_digest(_device_hwid_token(hwid), token):
target_hwid = hwid
break
if not target_hwid:
return _json_error(404, "device_not_found", "Device not found")
success = await panel_service.disconnect_device(panel_user_uuid, target_hwid)
if not success:
return _json_error(502, "device_disconnect_failed", "Failed to disconnect device")
await session.commit()
return web.json_response({"ok": True})
def _device_hwid_token(hwid: str) -> str:
return hashlib.sha256(str(hwid or "").encode()).hexdigest()[:32]
def _shorten_hwid_for_display(hwid: Optional[str], max_length: int = 24) -> str:
value = str(hwid or "").strip()
if len(value) <= max_length:
return value
return f"{value[:8]}...{value[-6:]}"
def _normalize_devices_response(devices_response: Any) -> List[Dict[str, Any]]:
if isinstance(devices_response, dict):
devices = devices_response.get("devices") or []
else:
devices = devices_response or []
if not isinstance(devices, list):
return []
return [device for device in devices if isinstance(device, dict)]
def _format_devices_limit(max_devices: Optional[int]) -> str:
if max_devices in (None, 0):
return "Unlimited"
return str(max_devices)
def _format_device_datetime(value: Any) -> str:
if not value:
return ""
text = str(value)
try:
normalized = datetime.fromisoformat(text.replace("Z", "+00:00"))
return normalized.strftime("%d.%m.%Y %H:%M")
except Exception:
return text
def _serialize_device(device: Dict[str, Any], index: int) -> Dict[str, Any]:
hwid = str(device.get("hwid") or "").strip()
model = str(device.get("deviceModel") or "").strip()
platform = str(device.get("platform") or "").strip()
os_version = str(device.get("osVersion") or "").strip()
user_agent = str(device.get("userAgent") or "").strip()
display_name = model or platform or f"Device {index}"
platform_label = " ".join(part for part in (platform, os_version) if part).strip()
return {
"index": index,
"display_name": display_name,
"platform": platform,
"os_version": os_version,
"platform_label": platform_label,
"user_agent": user_agent,
"created_at": device.get("createdAt"),
"created_at_text": _format_device_datetime(device.get("createdAt")),
"hwid_short": _shorten_hwid_for_display(hwid),
"token": _device_hwid_token(hwid) if hwid else "",
"can_disconnect": bool(hwid),
}
+59
View File
@@ -0,0 +1,59 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
class WebAppEmailPayload(BaseModel):
model_config = ConfigDict(extra="ignore")
email: EmailStr
@field_validator("email")
@classmethod
def _normalize_and_limit_email(cls, value: EmailStr) -> str:
normalized = normalize_email(str(value))
if len(normalized) > 254:
raise ValueError("email_too_long")
return normalized
class WebAppEmailCodePayload(WebAppEmailPayload):
code: str = ""
class WebAppEmailMagicPayload(BaseModel):
model_config = ConfigDict(extra="ignore")
token: constr(min_length=8, max_length=512)
class WebAppPaymentCreatePayload(BaseModel):
model_config = ConfigDict(extra="ignore")
method: str = ""
months: Any = None
traffic_gb: Any = None
device_count: Any = None
tariff_key: Optional[constr(max_length=128)] = None
sale_mode: Optional[constr(max_length=64)] = None
description: Optional[constr(max_length=4096)] = None
comment: Optional[constr(max_length=4096)] = None
note: Optional[constr(max_length=4096)] = None
class WebAppTariffChangePayload(BaseModel):
model_config = ConfigDict(extra="ignore")
tariff_key: constr(min_length=1, max_length=128)
mode: constr(min_length=1, max_length=64)
class WebAppLanguagePayload(BaseModel):
model_config = ConfigDict(extra="ignore")
language: constr(min_length=2, max_length=16)
class WebAppDeviceDisconnectPayload(BaseModel):
model_config = ConfigDict(extra="ignore")
token: constr(min_length=8, max_length=128)
+48
View File
@@ -0,0 +1,48 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
def setup_subscription_webapp_routes(app: web.Application) -> None:
app.router.add_get("/", index_route)
app.router.add_get("/home", index_route)
app.router.add_get("/invite", index_route)
app.router.add_get("/devices", index_route)
app.router.add_get("/settings", index_route)
app.router.add_get("/admin", index_route)
app.router.add_get("/admin/{section:[a-z][a-z0-9_-]*}", index_route)
app.router.add_get("/admin/users/{user_id:-?[0-9]+}", index_route)
app.router.add_get("/auth/telegram/start", telegram_oauth_start_route)
app.router.add_get("/auth/telegram/callback", telegram_oauth_callback_route)
app.router.add_get("/health", health_route)
app.router.add_get(WEBAPP_LOGO_PROXY_PATH, webapp_logo_route)
app.router.add_get(
r"/webapp-emoji/{codepoints:[0-9a-f_]+}/512.{ext:gif|webp}",
webapp_animated_emoji_route,
)
app.router.add_get("/subscription_webapp.css", css_asset_route)
app.router.add_get("/subscription_webapp.min.{asset_hash}.js", js_asset_route)
app.router.add_get("/subscription_webapp.js", js_asset_route)
app.router.add_post("/api/auth/telegram/nonce", telegram_oauth_nonce_route)
app.router.add_post("/api/auth/token", auth_token_route)
app.router.add_post("/api/auth/email/request", email_auth_request_route)
app.router.add_post("/api/auth/email/verify", email_auth_verify_route)
app.router.add_post("/api/auth/email/magic", email_auth_magic_route)
app.router.add_post("/api/auth/logout", logout_route)
app.router.add_get("/api/me", me_route)
app.router.add_get("/api/account/avatar", account_avatar_route)
app.router.add_post("/api/account/language", account_language_route)
app.router.add_post("/api/account/email/request", account_email_request_route)
app.router.add_post("/api/account/email/verify", account_email_verify_route)
app.router.add_post("/api/account/telegram/link", account_telegram_link_route)
app.router.add_post("/api/promo/apply", apply_promo_route)
app.router.add_post("/api/trial/activate", activate_trial_route)
app.router.add_get("/api/devices", devices_route)
app.router.add_post("/api/devices/disconnect", disconnect_device_route)
app.router.add_get("/api/devices/topup-options", device_topup_options_route)
app.router.add_get("/api/tariffs/topup-options", tariff_topup_options_route)
app.router.add_get("/api/tariffs/change-options", tariff_change_options_route)
app.router.add_post("/api/tariffs/change", tariff_change_route)
app.router.add_post("/api/tariffs/change-payment", tariff_change_payment_route)
app.router.add_post("/api/payments", create_payment_route)
app.router.add_get("/api/payments/{payment_id}", payment_status_route)
setup_admin_routes(app)
+612
View File
@@ -0,0 +1,612 @@
# ruff: noqa: F401,F403,F405,I001
from ._runtime import * # noqa: F403,F405
async def _build_user_payload(request: web.Request, user_id: int) -> Dict[str, Any]:
settings: Settings = request.app["settings"]
async_session_factory: sessionmaker = request.app["async_session_factory"]
subscription_service: SubscriptionService = request.app["subscription_service"]
cached = _get_cached_webapp_settings(request)
async with async_session_factory() as session:
db_user = await user_dal.get_user_by_id(session, user_id)
if not db_user or db_user.is_banned:
raise web.HTTPForbidden(
text=json.dumps({"ok": False, "error": "access_denied"}),
content_type="application/json",
)
active = await subscription_service.get_active_subscription_details(session, user_id)
referral_code = await user_dal.ensure_referral_code(session, db_user)
referral_service: Optional[ReferralService] = request.app.get("referral_service")
bot_username = request.app.get("bot_username") or ""
referral_link = None
if referral_service and bot_username:
referral_link = await referral_service.generate_referral_link(
session,
bot_username,
user_id,
)
webapp_referral_link = _build_webapp_referral_link(
request.app["settings"].SUBSCRIPTION_MINI_APP_URL,
referral_code,
)
referral_stats = (
await referral_service.get_referral_stats(session, user_id)
if referral_service
else {"invited_count": 0, "purchased_count": 0}
)
local_sub = (
await subscription_dal.get_active_subscription_by_user_id(
session,
user_id,
db_user.panel_user_uuid,
)
if db_user.panel_user_uuid
else None
)
trial_available = bool(
settings.TRIAL_ENABLED
and settings.TRIAL_DURATION_DAYS > 0
and not await subscription_service.has_had_any_subscription(session, user_id)
)
avatar = await _ensure_cached_telegram_avatar(request, session, db_user)
try:
await session.commit()
except Exception:
await session.rollback()
lang = _normalize_language(db_user.language_code or settings.DEFAULT_LANGUAGE)
admin_ids = {int(x) for x in (settings.ADMIN_IDS or [])}
is_admin = bool(db_user.telegram_id and int(db_user.telegram_id) in admin_ids)
return {
"user": {
"id": user_id,
"username": db_user.username,
"email": db_user.email,
"email_verified": bool(db_user.email_verified_at),
"telegram_id": db_user.telegram_id,
"telegram_linked": bool(_telegram_id_for_user(db_user)),
"telegram_photo_url": _telegram_avatar_url(avatar),
"first_name": db_user.first_name,
"language_code": lang,
"is_admin": is_admin,
},
"subscription": _serialize_subscription(settings, active, local_sub, lang),
"referral": {
"code": referral_code,
"bot_link": referral_link,
"webapp_link": webapp_referral_link,
"invited_count": referral_stats.get("invited_count", 0),
"purchased_count": referral_stats.get("purchased_count", 0),
"welcome_bonus_days": max(
0, int(getattr(settings, "REFERRAL_WELCOME_BONUS_DAYS", 0) or 0)
),
"one_bonus_per_referee": bool(
getattr(settings, "REFERRAL_ONE_BONUS_PER_REFEREE", False)
),
"bonus_details": _serialize_referral_bonus_details(settings, lang),
},
"plans": _serialize_plans(
settings,
lang,
subscription_options=cached["subscription_options"],
stars_subscription_options=cached["stars_subscription_options"],
traffic_packages=cached["traffic_packages"],
stars_traffic_packages=cached["stars_traffic_packages"],
),
"payment_methods": _serialize_payment_methods(settings, request.app),
"settings": {
"support_url": settings.SUPPORT_LINK,
"traffic_mode": bool(settings.traffic_sale_mode),
"my_devices_enabled": bool(settings.MY_DEVICES_SECTION_ENABLED),
"user_hwid_device_limit": (
int(settings.USER_HWID_DEVICE_LIMIT)
if settings.USER_HWID_DEVICE_LIMIT is not None
else None
),
"trial_enabled": bool(settings.TRIAL_ENABLED),
"trial_available": trial_available,
"trial_duration_days": int(settings.TRIAL_DURATION_DAYS or 0),
"trial_traffic_limit_gb": float(settings.TRIAL_TRAFFIC_LIMIT_GB or 0),
"trial_traffic_strategy": getattr(settings, "TRIAL_TRAFFIC_STRATEGY", "NO_RESET"),
"email_auth_enabled": settings.email_auth_configured,
},
}
def _serialize_referral_bonus_details(settings: Settings, lang: str) -> List[Dict[str, Any]]:
if getattr(settings, "traffic_sale_mode", False):
return []
details: List[Dict[str, Any]] = []
for months, _price in sorted(settings.subscription_options.items()):
inviter_days = settings.referral_bonus_inviter.get(months)
friend_days = settings.referral_bonus_referee.get(months)
if inviter_days is None and friend_days is None:
continue
details.append(
{
"months": int(months),
"title": _format_months_title(int(months), lang),
"inviter_days": int(inviter_days or 0),
"friend_days": int(friend_days or 0),
}
)
return details
def _build_webapp_referral_link(
base_url: Optional[str],
referral_code: Optional[str],
) -> Optional[str]:
if not base_url or not referral_code:
return None
parts = urlsplit(base_url)
query = dict(parse_qsl(parts.query, keep_blank_values=True))
query["ref"] = f"u{referral_code}"
return urlunsplit(
(
parts.scheme,
parts.netloc,
parts.path or "/",
urlencode(query),
parts.fragment,
)
)
def _serialize_subscription(
settings: Settings,
active: Optional[Dict[str, Any]],
local_sub: Optional[Any],
lang: str,
) -> Dict[str, Any]:
if not active:
return {
"active": False,
"status": "INACTIVE",
"remaining_text": _format_remaining(0, lang),
"days_left": 0,
"config_link": None,
"connect_url": None,
}
end_date = active.get("end_date")
if end_date and end_date.tzinfo is None:
end_date = end_date.replace(tzinfo=timezone.utc)
seconds_left = 0
if end_date:
seconds_left = max(
0,
int((end_date - datetime.now(timezone.utc)).total_seconds()),
)
can_topup_regular_traffic = False
can_topup_premium_traffic = False
can_topup_traffic = False
if settings.tariffs_config and active.get("tariff_key"):
try:
tariff = settings.tariffs_config.require(str(active.get("tariff_key")))
packages = settings.tariffs_config.topup_packages_for(tariff)
can_topup_regular_traffic = bool(packages and packages.has_any())
can_topup_premium_traffic = bool(
tariff.premium_squad_uuids
and tariff.premium_topup_packages
and tariff.premium_topup_packages.has_any()
)
can_topup_traffic = bool(can_topup_regular_traffic or can_topup_premium_traffic)
except Exception:
can_topup_regular_traffic = False
can_topup_premium_traffic = False
can_topup_traffic = False
return {
"active": seconds_left > 0,
"status": active.get("status_from_panel") or "UNKNOWN",
"end_date": end_date.isoformat() if end_date else None,
"end_date_text": end_date.strftime("%d.%m.%Y %H:%M") if end_date else "N/A",
"days_left": seconds_left // 86400,
"remaining_text": _format_remaining(seconds_left, lang),
"config_link": active.get("config_link"),
"connect_url": active.get("connect_button_url") or active.get("config_link"),
"traffic_limit": _format_bytes(active.get("traffic_limit_bytes"), zero_as_unlimited=True),
"traffic_used": _format_bytes(active.get("traffic_used_bytes")),
"traffic_limit_bytes": _coerce_int_or_none(active.get("traffic_limit_bytes")),
"traffic_used_bytes": _coerce_int_or_none(active.get("traffic_used_bytes")),
"tariff_key": active.get("tariff_key"),
"tariff_name": active.get("tariff_name"),
"tariff_description": active.get("tariff_description"),
"premium_title": active.get("premium_title"),
"billing_model": active.get("billing_model"),
"traffic_limit_strategy": str(active.get("traffic_limit_strategy") or ""),
"tier_baseline_bytes": _coerce_int_or_none(active.get("tier_baseline_bytes")),
"topup_balance_bytes": _coerce_int_or_none(active.get("topup_balance_bytes")),
"premium_limit": _format_bytes(active.get("premium_limit_bytes"), zero_as_unlimited=True),
"premium_used": _format_bytes(active.get("premium_used_bytes")),
"premium_limit_bytes": _coerce_int_or_none(active.get("premium_limit_bytes")),
"premium_used_bytes": _coerce_int_or_none(active.get("premium_used_bytes")),
"premium_baseline_bytes": _coerce_int_or_none(active.get("premium_baseline_bytes")),
"premium_topup_balance_bytes": _coerce_int_or_none(
active.get("premium_topup_balance_bytes")
),
"premium_topup_used_bytes": _coerce_int_or_none(active.get("premium_topup_used_bytes")),
"premium_bonus_bytes": _coerce_int_or_none(active.get("premium_bonus_bytes")) or 0,
"regular_bonus_bytes": _coerce_int_or_none(active.get("regular_bonus_bytes")) or 0,
"regular_unlimited_override": bool(active.get("regular_unlimited_override")),
"premium_unlimited_override": bool(active.get("premium_unlimited_override")),
"premium_is_limited": bool(active.get("premium_is_limited")),
"premium_squad_labels": list(active.get("premium_squad_labels") or []),
"premium_node_labels": list(active.get("premium_node_labels") or []),
"can_topup_traffic": can_topup_traffic,
"can_topup_regular_traffic": can_topup_regular_traffic,
"can_topup_premium_traffic": can_topup_premium_traffic,
"period_start_at": active.get("period_start_at").isoformat()
if active.get("period_start_at")
else None,
"is_throttled": bool(active.get("is_throttled")),
"max_devices": _coerce_int_or_none(active.get("max_devices")),
"base_hwid_device_limit": _coerce_int_or_none(active.get("base_hwid_device_limit")),
"extra_hwid_devices": _coerce_int_or_none(active.get("extra_hwid_devices")) or 0,
"auto_renew_enabled": bool(getattr(local_sub, "auto_renew_enabled", False)),
"provider": getattr(local_sub, "provider", None),
}
def _serialize_plans(
settings: Settings,
lang: str,
*,
subscription_options: Optional[Dict[int, float]] = None,
stars_subscription_options: Optional[Dict[int, int]] = None,
traffic_packages: Optional[Dict[float, float]] = None,
stars_traffic_packages: Optional[Dict[float, int]] = None,
) -> List[Dict[str, Any]]:
tariffs_config = settings.tariffs_config
if tariffs_config:
plans: List[Dict[str, Any]] = []
for tariff in tariffs_config.enabled_tariffs:
common = {
"tariff_key": tariff.key,
"tariff_name": tariff.name(lang),
"billing_model": tariff.billing_model,
"description": tariff.description(lang),
"squad_uuids": tariff.squad_uuids,
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"hwid_device_limit": tariff.hwid_device_limit,
"hwid_device_packages": _serialize_hwid_device_packages(
settings,
tariff,
tariff.hwid_device_packages,
lang,
),
}
if tariff.billing_model == "period":
for months in sorted(tariff.enabled_periods):
price = tariff.period_price(int(months), "rub")
stars_price = tariff.period_price(int(months), "stars")
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
plan = {
**common,
"id": f"{tariff.key}:period:{int(months)}",
"sale_mode": "subscription",
"months": int(months),
"price": float(price or 0),
"title": tariff.name(lang),
"subtitle": _format_months_title(int(months), lang),
"monthly_gb": tariff.monthly_gb,
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
else:
rub_packages = {
float(package.gb): float(package.price)
for package in (tariff.traffic_packages.rub if tariff.traffic_packages else [])
}
stars_packages = {
float(package.gb): int(float(package.price))
for package in (
tariff.traffic_packages.stars if tariff.traffic_packages else []
)
}
for traffic_gb in sorted(set(rub_packages) | set(stars_packages)):
price = rub_packages.get(traffic_gb)
stars_price = stars_packages.get(traffic_gb)
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
traffic_value = float(traffic_gb)
plan = {
**common,
"id": f"{tariff.key}:traffic:{_format_number_for_payload(traffic_value)}",
"sale_mode": "traffic_package",
"months": int(traffic_value)
if traffic_value.is_integer()
else traffic_value,
"traffic_gb": traffic_value,
"price": float(price or 0),
"title": tariff.name(lang),
"subtitle": _format_traffic_title(traffic_value, lang),
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
return plans
if getattr(settings, "traffic_sale_mode", False):
active_traffic_packages = traffic_packages or settings.traffic_packages
active_stars_traffic_packages = stars_traffic_packages or settings.stars_traffic_packages
traffic_units = sorted(set(active_traffic_packages) | set(active_stars_traffic_packages))
plans: List[Dict[str, Any]] = []
for traffic_gb in traffic_units:
price = active_traffic_packages.get(traffic_gb)
stars_price = active_stars_traffic_packages.get(traffic_gb)
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
traffic_value = float(traffic_gb)
plan = {
"months": int(traffic_value) if traffic_value.is_integer() else traffic_value,
"traffic_gb": traffic_value,
"price": float(price or 0),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"title": _format_traffic_title(traffic_value, lang),
"sale_mode": "traffic",
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
return plans
active_subscription_options = subscription_options or settings.subscription_options
active_stars_subscription_options = (
stars_subscription_options or settings.stars_subscription_options
)
plans: List[Dict[str, Any]] = []
for months in sorted(set(active_subscription_options) | set(active_stars_subscription_options)):
price = active_subscription_options.get(months)
stars_price = active_stars_subscription_options.get(months)
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
plan = {
"months": int(months),
"price": float(price or 0),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"title": _format_months_title(int(months), lang),
"sale_mode": "subscription",
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
return plans
def _traffic_percent(used: Optional[int], limit: Optional[int]) -> int:
used_val = int(used or 0)
limit_val = int(limit or 0)
if limit_val <= 0:
return 0
return max(0, min(100, round((used_val / limit_val) * 100)))
def _serialize_topup_packages(
settings: Settings,
tariff: Any,
packages: Optional[Any],
lang: str,
*,
sale_mode: str = "topup",
title_prefix: str = "",
) -> List[Dict[str, Any]]:
rub_packages = {
float(package.gb): float(package.price) for package in (packages.rub if packages else [])
}
stars_packages = {
float(package.gb): int(float(package.price))
for package in (packages.stars if packages else [])
}
plans: List[Dict[str, Any]] = []
for traffic_gb in sorted(set(rub_packages) | set(stars_packages)):
price = rub_packages.get(traffic_gb)
stars_price = stars_packages.get(traffic_gb)
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
traffic_value = float(traffic_gb)
plan: Dict[str, Any] = {
"id": f"{tariff.key}:{sale_mode}:{_format_number_for_payload(traffic_value)}",
"tariff_key": tariff.key,
"tariff_name": tariff.name(lang),
"billing_model": tariff.billing_model,
"sale_mode": sale_mode,
"months": int(traffic_value) if traffic_value.is_integer() else traffic_value,
"traffic_gb": traffic_value,
"price": float(price or 0),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"title": f"{title_prefix}{_format_traffic_title(traffic_value, lang)}",
"subtitle": tariff.premium_name(lang)
if sale_mode == "premium_topup"
else tariff.name(lang),
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
return plans
def _serialize_hwid_device_packages(
settings: Settings,
tariff: Any,
packages: Optional[Any],
lang: str,
) -> List[Dict[str, Any]]:
rub_packages = {
int(package.count): float(package.price) for package in (packages.rub if packages else [])
}
stars_packages = {
int(package.count): int(float(package.price))
for package in (packages.stars if packages else [])
}
plans: List[Dict[str, Any]] = []
for count in sorted(set(rub_packages) | set(stars_packages)):
price = rub_packages.get(count)
stars_price = stars_packages.get(count)
if price is None and (stars_price is None or int(stars_price) <= 0):
continue
plan: Dict[str, Any] = {
"id": f"{tariff.key}:hwid:{count}",
"tariff_key": tariff.key,
"tariff_name": tariff.name(lang),
"billing_model": tariff.billing_model,
"sale_mode": "hwid_devices",
"months": int(count),
"device_count": int(count),
"price": float(price or 0),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
"title": f"+{count}",
"subtitle": tariff.name(lang),
}
if stars_price is not None and int(stars_price) > 0:
plan["stars_price"] = int(stars_price)
plans.append(plan)
return plans
def _serialize_tariff_change_target(
settings: Settings,
config: Any,
tariff: Any,
options: Dict[str, Any],
lang: str,
) -> Dict[str, Any]:
actions: List[Dict[str, Any]] = []
mode = str(options.get("mode") or "")
if mode == "period_to_period":
actions.append(
{
"mode": "recalc_days",
"kind": "free",
"title": "recalc_days",
"days_after": int(options.get("recalc_days") or 0),
"remaining_days": int(options.get("remaining_days") or 0),
}
)
paid_diff = float(options.get("paid_diff_rub") or 0)
if paid_diff > 0:
actions.append(
{
"mode": "paid_diff",
"kind": "payment",
"title": "paid_diff",
"price": paid_diff,
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
}
)
elif mode == "period_to_traffic":
actions.append(
{
"mode": "convert_days_to_gb",
"kind": "free",
"title": "convert_days_to_gb",
"converted_gb": float(options.get("converted_gb") or 0),
"remaining_days": int(options.get("remaining_days") or 0),
}
)
actions.extend(
{
"mode": "buy_package",
"kind": "payment",
"title": f"+{package.gb:g} GB",
"traffic_gb": float(package.gb),
"price": float(package.price),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
}
for package in (tariff.traffic_packages.rub if tariff.traffic_packages else [])
)
else:
for months in tariff.enabled_periods:
price = tariff.period_price(int(months), "rub")
if price:
actions.append(
{
"mode": "buy_period",
"kind": "payment",
"months": int(months),
"title": _format_months_title(int(months), lang),
"price": float(price),
"currency": settings.DEFAULT_CURRENCY_SYMBOL or "RUB",
}
)
return {
"tariff_key": tariff.key,
"title": tariff.name(lang),
"description": tariff.description(lang),
"billing_model": tariff.billing_model,
"monthly_gb": tariff.monthly_gb,
"options": options,
"actions": actions,
}
def _serialize_payment_methods(
settings: Settings,
app: web.Application,
) -> List[Dict[str, Any]]:
labels = {
"severpay": "SeverPay",
"freekassa": "FreeKassa / СБП",
"platega_sbp": "Platega · СБП",
"platega_crypto": "Platega · Crypto",
"yookassa": "Банковская карта",
"stars": "Telegram Stars",
"cryptopay": "CryptoPay",
}
methods: List[Dict[str, Any]] = []
for method in settings.payment_methods_order:
method = method.lower()
if (
method == "severpay"
and settings.SEVERPAY_ENABLED
and _service_configured(app, "severpay_service")
):
methods.append({"id": method, "name": labels[method]})
elif (
method == "freekassa"
and settings.FREEKASSA_ENABLED
and _service_configured(app, "freekassa_service")
):
methods.append({"id": method, "name": labels[method]})
elif (
method == "platega_sbp"
and settings.PLATEGA_ENABLED
and settings.PLATEGA_SBP_ENABLED
and _service_configured(app, "platega_service")
):
methods.append({"id": method, "name": labels[method]})
elif (
method == "platega_crypto"
and settings.PLATEGA_ENABLED
and settings.PLATEGA_CRYPTO_ENABLED
and _service_configured(app, "platega_service")
):
methods.append({"id": method, "name": labels[method]})
elif (
method == "yookassa"
and settings.YOOKASSA_ENABLED
and _service_configured(app, "yookassa_service")
):
methods.append({"id": method, "name": labels[method]})
elif method == "stars" and settings.STARS_ENABLED:
methods.append({"id": method, "name": labels[method]})
elif (
method == "cryptopay"
and settings.CRYPTOPAY_ENABLED
and _service_configured(app, "cryptopay_service")
):
methods.append({"id": method, "name": labels[method]})
return methods
def _service_configured(app: web.Application, key: str) -> bool:
service = app.get(key)
return bool(service and getattr(service, "configured", False))