diff --git a/.env.example b/.env.example index 37366091..f1d6c807 100644 --- a/.env.example +++ b/.env.example @@ -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 # === Ротация логов === # Включить новую систему ротации (по умолчанию старое поведение) diff --git a/app/cabinet/routes/admin_campaigns.py b/app/cabinet/routes/admin_campaigns.py index 5e922b1e..38f4289e 100644 --- a/app/cabinet/routes/admin_campaigns.py +++ b/app/cabinet/routes/admin_campaigns.py @@ -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, ) ) diff --git a/app/cabinet/routes/auth.py b/app/cabinet/routes/auth.py index 20c40309..5a3a64b4 100644 --- a/app/cabinet/routes/auth.py +++ b/app/cabinet/routes/auth.py @@ -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 diff --git a/app/cabinet/routes/oauth.py b/app/cabinet/routes/oauth.py index 9142a3df..90863801 100644 --- a/app/cabinet/routes/oauth.py +++ b/app/cabinet/routes/oauth.py @@ -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) diff --git a/app/cabinet/schemas/auth.py b/app/cabinet/schemas/auth.py index c92beb12..58f31340 100644 --- a/app/cabinet/schemas/auth.py +++ b/app/cabinet/schemas/auth.py @@ -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): diff --git a/app/cabinet/schemas/campaigns.py b/app/cabinet/schemas/campaigns.py index c0c0530d..d5f98f2d 100644 --- a/app/cabinet/schemas/campaigns.py +++ b/app/cabinet/schemas/campaigns.py @@ -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): diff --git a/app/config.py b/app/config.py index 31dc929c..826e005d 100644 --- a/app/config.py +++ b/app/config.py @@ -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 # По умолчанию старое поведение diff --git a/app/database/crud/campaign.py b/app/database/crud/campaign.py index 2c3717b8..8b8f2f37 100644 --- a/app/database/crud/campaign.py +++ b/app/database/crud/campaign.py @@ -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) diff --git a/app/handlers/referral.py b/app/handlers/referral.py index 50701f7c..7f667cfa 100644 --- a/app/handlers/referral.py +++ b/app/handlers/referral.py @@ -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) diff --git a/app/handlers/subscription/purchase.py b/app/handlers/subscription/purchase.py index 1247b1f8..cd0c82e3 100644 --- a/app/handlers/subscription/purchase.py +++ b/app/handlers/subscription/purchase.py @@ -348,8 +348,6 @@ async def show_subscription_info(callback: types.CallbackQuery, db_user: User, d tariff_info_lines.append('⏸️ Подписка приостановлена') # Показываем оставшееся время даже при паузе 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) diff --git a/app/logging_config.py b/app/logging_config.py index b1308c98..3f189749 100644 --- a/app/logging_config.py +++ b/app/logging_config.py @@ -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), ], ) diff --git a/app/services/admin_notification_service.py b/app/services/admin_notification_service.py index c975b171..7b13d5f1 100644 --- a/app/services/admin_notification_service.py +++ b/app/services/admin_notification_service.py @@ -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. diff --git a/app/services/campaign_service.py b/app/services/campaign_service.py index 68db2c48..87ea2305 100644 --- a/app/services/campaign_service.py +++ b/app/services/campaign_service.py @@ -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,