Merge pull request #2618 from BEDOLAGA-DEV/dev

Dev
This commit is contained in:
Egor
2026-02-17 07:01:07 +03:00
committed by GitHub
13 changed files with 171 additions and 103 deletions
+2
View File
@@ -845,6 +845,8 @@ VERSION_CHECK_INTERVAL_HOURS=1
# ===== ЛОГИРОВАНИЕ =====
LOG_LEVEL=INFO
LOG_FILE=logs/bot.log
# ANSI-цвета в консоли (true — цветной вывод с Rich, false — plain-text)
LOG_COLORS=true
# === Ротация логов ===
# Включить новую систему ротации (по умолчанию старое поведение)
+21 -8
View File
@@ -57,6 +57,14 @@ def _get_deep_link(start_parameter: str) -> str:
return f'?start={start_parameter}'
def _get_web_link(start_parameter: str) -> str | None:
"""Generate web link for campaign."""
base_url = (settings.MINIAPP_CUSTOM_URL or '').rstrip('/')
if base_url:
return f'{base_url}/?campaign={start_parameter}'
return None
@router.get('/overview', response_model=CampaignsOverviewResponse)
async def get_overview(
admin: User = Depends(get_current_admin_user),
@@ -205,6 +213,7 @@ async def get_campaign(
created_at=campaign.created_at,
updated_at=campaign.updated_at,
deep_link=_get_deep_link(campaign.start_parameter),
web_link=_get_web_link(campaign.start_parameter),
)
@@ -248,6 +257,7 @@ async def get_campaign_stats(
conversion_rate=stats['conversion_rate'],
trial_conversion_rate=stats['trial_conversion_rate'],
deep_link=_get_deep_link(campaign.start_parameter),
web_link=_get_web_link(campaign.start_parameter),
)
@@ -288,19 +298,22 @@ async def get_campaign_registrations(
)
total = count_result.scalar() or 0
items = []
for reg, user in rows:
# Check if user has subscription
# Batch query: find which users have active subscriptions (avoids N+1)
user_ids = [user.id for _reg, user in rows]
active_sub_user_ids: set[int] = set()
if user_ids:
sub_result = await db.execute(
select(Subscription)
select(Subscription.user_id)
.where(
Subscription.user_id == user.id,
Subscription.user_id.in_(user_ids),
Subscription.status == 'active',
)
.limit(1)
.distinct()
)
has_sub = sub_result.scalar_one_or_none() is not None
active_sub_user_ids = set(sub_result.scalars().all())
items = []
for reg, user in rows:
items.append(
CampaignRegistrationItem(
id=reg.id,
@@ -315,7 +328,7 @@ async def get_campaign_registrations(
tariff_duration_days=reg.tariff_duration_days,
created_at=reg.created_at,
user_balance_kopeks=user.balance_kopeks or 0,
has_subscription=has_sub,
has_subscription=user.id in active_sub_user_ids,
has_paid=user.has_had_paid_subscription or False,
)
)
+76 -9
View File
@@ -11,6 +11,10 @@ from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.database.crud.campaign import (
get_campaign_by_start_parameter,
get_campaign_registration_by_user,
)
from app.database.crud.user import (
clear_email_change_pending,
create_user,
@@ -23,6 +27,7 @@ from app.database.crud.user import (
verify_and_apply_email_change,
)
from app.database.models import CabinetRefreshToken, User
from app.services.campaign_service import AdvertisingCampaignService
from app.services.disposable_email_service import disposable_email_service
from app.services.referral_service import process_referral_registration
from app.utils.timezone import panel_datetime_to_utc
@@ -49,6 +54,7 @@ from ..auth.jwt_handler import get_refresh_token_expires_at
from ..dependencies import get_cabinet_db, get_current_cabinet_user
from ..schemas.auth import (
AuthResponse,
CampaignBonusInfo,
EmailChangeRequest,
EmailChangeResponse,
EmailChangeVerifyRequest,
@@ -118,12 +124,6 @@ async def _store_refresh_token(
token_hash = hashlib.sha256(refresh_token.encode()).hexdigest()
expires_at = get_refresh_token_expires_at()
# Check if token already exists (handles race conditions)
existing = await db.execute(select(CabinetRefreshToken).where(CabinetRefreshToken.token_hash == token_hash))
if existing.scalar_one_or_none():
# Token already stored, skip
return
token_record = CabinetRefreshToken(
user_id=user_id,
token_hash=token_hash,
@@ -133,9 +133,56 @@ async def _store_refresh_token(
db.add(token_record)
try:
await db.commit()
except Exception:
# Handle race condition if token was inserted between check and insert
except IntegrityError:
await db.rollback()
logger.debug('Refresh token already exists (duplicate)', user_id=user_id)
async def _process_campaign_bonus(
db: AsyncSession,
user: User,
campaign_slug: str | None,
) -> CampaignBonusInfo | None:
"""Process campaign bonus for user during auth. Never raises."""
if not campaign_slug:
return None
try:
campaign = await get_campaign_by_start_parameter(db, campaign_slug, only_active=True)
if not campaign:
return None
# Lock user row to prevent concurrent bonus application (race condition)
await db.execute(select(User).where(User.id == user.id).with_for_update())
existing = await get_campaign_registration_by_user(db, user.id)
if existing:
logger.debug('User already has campaign registration', user_id=user.id)
return None
service = AdvertisingCampaignService()
result = await service.apply_campaign_bonus(db, user, campaign)
if not result.success:
return None
# Refresh user to get updated balance after bonus
await db.refresh(user)
return CampaignBonusInfo(
campaign_name=campaign.name,
bonus_type=result.bonus_type or campaign.bonus_type,
balance_kopeks=result.balance_kopeks,
subscription_days=result.subscription_days,
tariff_name=result.tariff_name,
)
except Exception:
logger.exception('Failed to process campaign bonus', user_id=user.id, campaign_slug=campaign_slug)
try:
await db.rollback()
# Re-fetch user so session stays usable for the caller
await db.refresh(user)
except Exception:
logger.exception('Failed to rollback after campaign bonus error', user_id=user.id)
return None
async def _sync_subscription_from_panel_by_email(db: AsyncSession, user: User) -> None:
@@ -321,6 +368,11 @@ async def auth_telegram(
# Store refresh token
await _store_refresh_token(db, user.id, response.refresh_token)
# Process campaign bonus
response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug)
if response.campaign_bonus:
response.user = _user_to_response(user)
return response
@@ -335,7 +387,7 @@ async def auth_telegram_widget(
This endpoint validates data from Telegram Login Widget and returns
JWT tokens for authenticated access.
"""
widget_data = request.model_dump()
widget_data = request.model_dump(exclude={'campaign_slug'})
if not validate_telegram_login_widget(widget_data):
raise HTTPException(
@@ -380,6 +432,11 @@ async def auth_telegram_widget(
response = _create_auth_response(user)
await _store_refresh_token(db, user.id, response.refresh_token)
# Process campaign bonus
response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug)
if response.campaign_bonus:
response.user = _user_to_response(user)
return response
@@ -647,6 +704,11 @@ async def verify_email(
response = _create_auth_response(user)
await _store_refresh_token(db, user.id, response.refresh_token)
# Process campaign bonus
response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug)
if response.campaign_bonus:
response.user = _user_to_response(user)
return response
@@ -789,6 +851,11 @@ async def login_email(
response = _create_auth_response(user)
await _store_refresh_token(db, user.id, response.refresh_token)
# Process campaign bonus
response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug)
if response.campaign_bonus:
response.user = _user_to_response(user)
return response
+18 -5
View File
@@ -24,7 +24,7 @@ from ..auth.oauth_providers import (
)
from ..dependencies import get_cabinet_db
from ..schemas.auth import AuthResponse
from .auth import _create_auth_response, _store_refresh_token
from .auth import _create_auth_response, _process_campaign_bonus, _store_refresh_token
logger = structlog.get_logger(__name__)
@@ -32,12 +32,22 @@ logger = structlog.get_logger(__name__)
router = APIRouter(prefix='/auth/oauth', tags=['Cabinet OAuth'])
async def _finalize_oauth_login(db: AsyncSession, user: User, provider: str) -> AuthResponse:
async def _finalize_oauth_login(
db: AsyncSession,
user: User,
provider: str,
campaign_slug: str | None = None,
) -> AuthResponse:
"""Update last login, create tokens, store refresh token."""
user.cabinet_last_login = datetime.now(UTC)
await db.commit()
auth_response = _create_auth_response(user)
await _store_refresh_token(db, user.id, auth_response.refresh_token, device_info=f'oauth:{provider}')
auth_response.campaign_bonus = await _process_campaign_bonus(db, user, campaign_slug)
if auth_response.campaign_bonus:
from .auth import _user_to_response
auth_response.user = _user_to_response(user)
return auth_response
@@ -61,6 +71,9 @@ class OAuthAuthorizeResponse(BaseModel):
class OAuthCallbackRequest(BaseModel):
code: str = Field(..., description='Authorization code from provider')
state: str = Field(..., description='CSRF state token')
campaign_slug: str | None = Field(
None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$', description='Campaign slug from web link'
)
# --- Endpoints ---
@@ -140,7 +153,7 @@ async def oauth_callback(
user = await get_user_by_oauth_provider(db, provider, user_info.provider_id)
if user:
logger.info('OAuth login via for existing user', provider=provider, user_id=user.id)
return await _finalize_oauth_login(db, user, provider)
return await _finalize_oauth_login(db, user, provider, request.campaign_slug)
# 6. Find user by email (if verified) and link provider
if user_info.email and user_info.email_verified:
@@ -148,7 +161,7 @@ async def oauth_callback(
if user:
await set_user_oauth_provider_id(db, user, provider, user_info.provider_id)
logger.info('OAuth login via linked to existing email user', provider=provider, user_id=user.id)
return await _finalize_oauth_login(db, user, provider)
return await _finalize_oauth_login(db, user, provider, request.campaign_slug)
# 7. Create new user
user = await create_user_by_oauth(
@@ -162,4 +175,4 @@ async def oauth_callback(
username=user_info.username,
)
logger.info('OAuth new user created via with id', provider=provider, user_id=user.id)
return await _finalize_oauth_login(db, user, provider)
return await _finalize_oauth_login(db, user, provider, request.campaign_slug)
+23
View File
@@ -9,6 +9,9 @@ class TelegramAuthRequest(BaseModel):
"""Request for Telegram WebApp initData authentication."""
init_data: str = Field(..., description='Telegram WebApp initData string')
campaign_slug: str | None = Field(
None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$', description='Campaign slug from web link'
)
class TelegramWidgetAuthRequest(BaseModel):
@@ -21,6 +24,9 @@ class TelegramWidgetAuthRequest(BaseModel):
photo_url: str | None = Field(None, description="User's photo URL")
auth_date: int = Field(..., description='Unix timestamp of authentication')
hash: str = Field(..., description='Authentication hash')
campaign_slug: str | None = Field(
None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$', description='Campaign slug from web link'
)
class EmailRegisterRequest(BaseModel):
@@ -34,6 +40,9 @@ class EmailVerifyRequest(BaseModel):
"""Request to verify email with token."""
token: str = Field(..., description='Email verification token')
campaign_slug: str | None = Field(
None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$', description='Campaign slug from web link'
)
class EmailLoginRequest(BaseModel):
@@ -41,6 +50,9 @@ class EmailLoginRequest(BaseModel):
email: EmailStr = Field(..., description='Email address')
password: str = Field(..., description='Password')
campaign_slug: str | None = Field(
None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$', description='Campaign slug from web link'
)
class RefreshTokenRequest(BaseModel):
@@ -102,6 +114,16 @@ class EmailRegisterStandaloneRequest(BaseModel):
referral_code: str | None = Field(None, max_length=32, description='Referral code of inviter')
class CampaignBonusInfo(BaseModel):
"""Info about campaign bonus applied during auth."""
campaign_name: str
bonus_type: str
balance_kopeks: int = 0
subscription_days: int | None = None
tariff_name: str | None = None
class AuthResponse(BaseModel):
"""Full authentication response with tokens and user."""
@@ -110,6 +132,7 @@ class AuthResponse(BaseModel):
token_type: str = 'bearer'
expires_in: int
user: UserResponse
campaign_bonus: CampaignBonusInfo | None = None
class RegisterResponse(BaseModel):
+4 -2
View File
@@ -66,6 +66,7 @@ class CampaignDetailResponse(BaseModel):
updated_at: datetime | None = None
# Deep link
deep_link: str | None = None
web_link: str | None = None
class Config:
from_attributes = True
@@ -75,7 +76,7 @@ class CampaignCreateRequest(BaseModel):
"""Request to create a campaign."""
name: str = Field(..., min_length=1, max_length=255)
start_parameter: str = Field(..., min_length=1, max_length=100, pattern=r'^[a-zA-Z0-9_-]+$')
start_parameter: str = Field(..., min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$')
bonus_type: CampaignBonusType
is_active: bool = True
# Balance bonus
@@ -94,7 +95,7 @@ class CampaignUpdateRequest(BaseModel):
"""Request to update a campaign."""
name: str | None = Field(None, min_length=1, max_length=255)
start_parameter: str | None = Field(None, min_length=1, max_length=100, pattern=r'^[a-zA-Z0-9_-]+$')
start_parameter: str | None = Field(None, min_length=1, max_length=64, pattern=r'^[a-zA-Z0-9_-]+$')
bonus_type: CampaignBonusType | None = None
is_active: bool | None = None
# Balance bonus
@@ -147,6 +148,7 @@ class CampaignStatisticsResponse(BaseModel):
trial_conversion_rate: float = 0.0
# Deep link
deep_link: str | None = None
web_link: str | None = None
class CampaignRegistrationItem(BaseModel):
+1
View File
@@ -549,6 +549,7 @@ class Settings(BaseSettings):
LOG_LEVEL: str = 'INFO'
LOG_FILE: str = 'logs/bot.log'
LOG_COLORS: bool = True # ANSI-цвета в консоли (false для plain-text вывода)
# === Log Rotation Settings ===
LOG_ROTATION_ENABLED: bool = False # По умолчанию старое поведение
-60
View File
@@ -359,66 +359,6 @@ async def get_campaign_statistics(
if first_payment_amount_by_user:
avg_first_payment = int(sum(first_payment_amount_by_user.values()) / len(first_payment_amount_by_user))
conversion_rate = 0.0
if count:
conversion_rate = round((paid_users_count / count) * 100, 1)
trial_conversion_rate = 0.0
if trial_users_count:
trial_conversion_rate = round((conversion_count / trial_users_count) * 100, 1)
avg_revenue_per_user = 0
if count:
avg_revenue_per_user = int(total_revenue / count)
deposits_result = await db.execute(
select(func.coalesce(func.sum(Transaction.amount_kopeks), 0)).where(
Transaction.user_id.in_(select(registrations_subquery.c.user_id)),
Transaction.type == TransactionType.DEPOSIT.value,
Transaction.is_completed.is_(True),
)
)
total_revenue = deposits_result.scalar() or 0
trials_result = await db.execute(
select(func.count(func.distinct(Subscription.user_id))).where(
Subscription.user_id.in_(select(registrations_subquery.c.user_id)),
Subscription.is_trial.is_(True),
)
)
trial_users_count = trials_result.scalar() or 0
active_trials_result = await db.execute(
select(func.count(func.distinct(Subscription.user_id))).where(
Subscription.user_id.in_(select(registrations_subquery.c.user_id)),
Subscription.is_trial.is_(True),
Subscription.status == SubscriptionStatus.ACTIVE.value,
)
)
active_trials_count = active_trials_result.scalar() or 0
conversions_result = await db.execute(
select(func.count(func.distinct(SubscriptionConversion.user_id))).where(
SubscriptionConversion.user_id.in_(select(registrations_subquery.c.user_id))
)
)
conversion_count = conversions_result.scalar() or 0
paid_users_result = await db.execute(
select(func.count(User.id)).where(
User.id.in_(select(registrations_subquery.c.user_id)),
User.has_had_paid_subscription.is_(True),
)
)
paid_users_count = paid_users_result.scalar() or 0
avg_first_payment_result = await db.execute(
select(func.coalesce(func.avg(SubscriptionConversion.first_payment_amount_kopeks), 0)).where(
SubscriptionConversion.user_id.in_(select(registrations_subquery.c.user_id))
)
)
avg_first_payment = int(avg_first_payment_result.scalar() or 0)
conversion_rate = 0.0
if count:
conversion_rate = round((paid_users_count / count) * 100, 1)
+1 -1
View File
@@ -793,7 +793,7 @@ async def confirm_withdrawal_request(callback: types.CallbackQuery, db_user: Use
try:
notification_service = AdminNotificationService(callback.bot)
await notification_service.send_to_admins(admin_text, keyboard=admin_keyboard)
await notification_service.send_admin_notification(admin_text, reply_markup=admin_keyboard)
except Exception as e:
logger.error('Ошибка отправки уведомления админам о заявке на вывод', error=e)
-4
View File
@@ -348,8 +348,6 @@ async def show_subscription_info(callback: types.CallbackQuery, db_user: User, d
tariff_info_lines.append('⏸️ <b>Подписка приостановлена</b>')
# Показываем оставшееся время даже при паузе
if last_charge:
from datetime import UTC, timedelta
next_charge = last_charge + timedelta(hours=24)
now = datetime.now(UTC)
if next_charge > now:
@@ -359,8 +357,6 @@ async def show_subscription_info(callback: types.CallbackQuery, db_user: User, d
tariff_info_lines.append(f'⏳ Осталось: {hours_left}ч {minutes_left}мин')
tariff_info_lines.append('💤 Списание приостановлено')
elif last_charge:
from datetime import UTC, timedelta
next_charge = last_charge + timedelta(hours=24)
now = datetime.now(UTC)
+19 -12
View File
@@ -121,24 +121,31 @@ def setup_logging() -> tuple[logging.Formatter, logging.Formatter, Any]:
],
)
# Console formatter: colors enabled by default on non-Windows.
# Console formatter: colors controlled by LOG_COLORS env var (default: true).
# Rich tracebacks with conservative limits to avoid 5000-line dumps.
use_colors = settings.LOG_COLORS
console_renderer_kwargs: dict[str, Any] = {
'colors': use_colors,
'pad_event_to': 0,
'pad_level': False,
}
if use_colors:
console_renderer_kwargs['exception_formatter'] = structlog.dev.RichTracebackFormatter(
show_locals=False,
max_frames=20,
extra_lines=1,
width=120,
suppress=['aiogram', 'aiohttp'],
)
else:
console_renderer_kwargs['exception_formatter'] = structlog.dev.plain_traceback
console_formatter = structlog.stdlib.ProcessorFormatter(
foreign_pre_chain=shared_processors,
processors=[
structlog.stdlib.ProcessorFormatter.remove_processors_meta,
_prefix_logger_name,
structlog.dev.ConsoleRenderer(
pad_event_to=0,
pad_level=False,
exception_formatter=structlog.dev.RichTracebackFormatter(
show_locals=False,
max_frames=20,
extra_lines=1,
width=120,
suppress=['aiogram', 'aiohttp'],
),
),
structlog.dev.ConsoleRenderer(**console_renderer_kwargs),
],
)
@@ -1215,6 +1215,12 @@ class AdminNotificationService:
"""Public check for whether admin notifications are configured and active."""
return self._is_enabled()
async def send_admin_notification(self, text: str, reply_markup: types.InlineKeyboardMarkup | None = None) -> bool:
"""Send a generic notification to admin chat with optional inline keyboard."""
if not self._is_enabled():
return False
return await self._send_message(text, reply_markup=reply_markup)
async def send_webhook_notification(self, text: str) -> bool:
"""Send a generic webhook/infrastructure notification to admin chat.
-2
View File
@@ -150,8 +150,6 @@ class AdvertisingCampaignService:
except Exception as error:
logger.error('Не удалось подобрать сквад для кампании', campaign_id=campaign.id, error=error)
squads[0] if squads else None
new_subscription = await create_paid_subscription(
db=db,
user_id=user.id,