diff --git a/bot/handlers/admin/__init__.py b/bot/handlers/admin/__init__.py index ecf8b9f..6f5dc6d 100644 --- a/bot/handlers/admin/__init__.py +++ b/bot/handlers/admin/__init__.py @@ -8,6 +8,7 @@ from . import statistics from . import sync_admin from . import logs_admin from . import payments +from . import ads admin_router_aggregate = Router(name="admin_features_router") @@ -19,5 +20,6 @@ admin_router_aggregate.include_router(statistics.router) admin_router_aggregate.include_router(sync_admin.router) admin_router_aggregate.include_router(logs_admin.router) admin_router_aggregate.include_router(payments.router) +admin_router_aggregate.include_router(ads.router) __all__ = ("admin_router_aggregate", ) diff --git a/bot/handlers/admin/ads.py b/bot/handlers/admin/ads.py new file mode 100644 index 0000000..88f5b1c --- /dev/null +++ b/bot/handlers/admin/ads.py @@ -0,0 +1,163 @@ +import logging +from aiogram import Router, F, types +from aiogram.fsm.context import FSMContext +from typing import Optional +from sqlalchemy.ext.asyncio import AsyncSession + +from config.settings import Settings +from bot.middlewares.i18n import JsonI18n +from db.dal import ad_dal + +router = Router(name="admin_ads_router") + + +@router.callback_query(F.data == "admin_action:ads") +async def show_ads_menu(callback: types.CallbackQuery, settings: Settings, i18n_data: dict, session: AsyncSession): + current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE) + i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") + _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs) if i18n else key + + if not i18n or not callback.message: + await callback.answer("Language error.", show_alert=True) + return + + campaigns = await ad_dal.list_campaigns(session) + if not campaigns: + text = _("admin_ads_empty") + else: + text_lines = [_("admin_ads_header")] + for camp in campaigns: + try: + stats = await ad_dal.get_campaign_stats(session, camp.ad_campaign_id) + except Exception as e_stats: + logging.error(f"Failed to calc stats for campaign {camp.ad_campaign_id}: {e_stats}") + stats = {"starts": 0, "trials": 0, "payers": 0, "revenue": 0.0} + text_lines.append( + _( + "admin_ads_item", + id=camp.ad_campaign_id, + source=camp.source, + start_param=camp.start_param, + cost=f"{camp.cost:.2f}", + active=_("csv_yes") if camp.is_active else _("csv_no"), + starts=stats["starts"], + trials=stats["trials"], + payers=stats["payers"], + revenue=f"{stats['revenue']:.2f}", + ) + ) + text = "\n\n".join(text_lines) + + from bot.keyboards.inline.admin_keyboards import get_ads_menu_keyboard + reply_markup = get_ads_menu_keyboard(i18n, current_lang) + await callback.message.edit_text(text, reply_markup=reply_markup) + try: + await callback.answer() + except Exception: + pass + + +@router.callback_query(F.data == "admin_action:ads_create") +async def ads_create_start(callback: types.CallbackQuery, state: FSMContext, settings: Settings, i18n_data: dict): + from bot.states.admin_states import AdminStates + current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE) + i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") + _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs) if i18n else key + + if not i18n or not callback.message: + await callback.answer("Language error.", show_alert=True) + return + + await state.set_state(AdminStates.waiting_for_ad_source) + await callback.message.edit_text(_("admin_ads_create_source_prompt")) + try: + await callback.answer() + except Exception: + pass + + +@router.message(F.text, state="*") +async def ads_create_flow(message: types.Message, state: FSMContext, settings: Settings, i18n_data: dict, session: AsyncSession): + from bot.states.admin_states import AdminStates + current_state = await state.get_state() + if current_state not in ( + AdminStates.waiting_for_ad_source.state, + AdminStates.waiting_for_ad_start_param.state, + AdminStates.waiting_for_ad_cost.state, + ): + return + + current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE) + i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") + _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs) if i18n else key + + if current_state == AdminStates.waiting_for_ad_source.state: + source = message.text.strip() + if not source or len(source) > 64: + await message.answer(_("admin_ads_invalid_source")) + return + await state.update_data(ad_source=source) + await state.set_state(AdminStates.waiting_for_ad_start_param) + await message.answer(_("admin_ads_create_start_param_prompt")) + return + + if current_state == AdminStates.waiting_for_ad_start_param.state: + start_param = message.text.strip() + # Allow alnum underscore dash only + import re as _re + if not _re.match(r"^[A-Za-z0-9_\-]{2,64}$", start_param): + await message.answer(_("admin_ads_invalid_start_param")) + return + await state.update_data(ad_start_param=start_param) + await state.set_state(AdminStates.waiting_for_ad_cost) + await message.answer(_("admin_ads_create_cost_prompt")) + return + + if current_state == AdminStates.waiting_for_ad_cost.state: + text = message.text.replace(",", ".").strip() + try: + cost = float(text) + if cost < 0 or cost > 1e8: + raise ValueError() + except Exception: + await message.answer(_("admin_ads_invalid_cost")) + return + + data = await state.get_data() + try: + campaign = await ad_dal.create_campaign( + session, + source=data.get("ad_source", "unknown"), + start_param=data.get("ad_start_param", "NA"), + cost=cost, + ) + await session.commit() + except ValueError as ve: + await session.rollback() + if str(ve) == "ad_campaign_start_param_exists": + await message.answer(_("admin_ads_start_param_exists")) + else: + await message.answer(_("error_occurred_try_again")) + return + except Exception as e: + await session.rollback() + logging.error(f"Failed to create ad campaign: {e}", exc_info=True) + await message.answer(_("error_occurred_try_again")) + return + + await state.clear() + _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs) if i18n else key + await message.answer( + _( + "admin_ads_created_success", + id=campaign.ad_campaign_id, + source=campaign.source, + start_param=campaign.start_param, + cost=f"{campaign.cost:.2f}", + ) + ) + # Offer back to ads menu + from bot.keyboards.inline.admin_keyboards import get_ads_menu_keyboard + await message.answer(_("admin_ads_back_to_menu_hint"), reply_markup=get_ads_menu_keyboard(i18n, current_lang)) + + diff --git a/bot/handlers/user/start.py b/bot/handlers/user/start.py index 0c25f17..294b055 100644 --- a/bot/handlers/user/start.py +++ b/bot/handlers/user/start.py @@ -118,6 +118,7 @@ async def send_main_menu(target_event: Union[types.Message, @router.message(CommandStart()) @router.message(CommandStart(magic=F.args.regexp(r"^ref_(\d+)$").as_("ref_match"))) @router.message(CommandStart(magic=F.args.regexp(r"^promo_(\w+)$").as_("promo_match"))) +@router.message(CommandStart(magic=F.args.regexp(r"^(?!ref_|promo_)([A-Za-z0-9_\-]{2,64})$").as_("ad_param_match"))) async def start_command_handler(message: types.Message, state: FSMContext, settings: Settings, @@ -125,7 +126,8 @@ async def start_command_handler(message: types.Message, subscription_service: SubscriptionService, session: AsyncSession, ref_match: Optional[re.Match] = None, - promo_match: Optional[re.Match] = None): + promo_match: Optional[re.Match] = None, + ad_param_match: Optional[re.Match] = None): await state.clear() current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE) i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") @@ -137,6 +139,7 @@ async def start_command_handler(message: types.Message, referred_by_user_id: Optional[int] = None promo_code_to_apply: Optional[str] = None + ad_start_param: Optional[str] = None if ref_match: potential_referrer_id = int(ref_match.group(1)) @@ -145,6 +148,9 @@ async def start_command_handler(message: types.Message, elif promo_match: promo_code_to_apply = promo_match.group(1) logging.info(f"User {user_id} started with promo code: {promo_code_to_apply}") + elif ad_param_match: + ad_start_param = ad_param_match.group(1) + logging.info(f"User {user_id} started with ad start param: {ad_start_param}") db_user = await user_dal.get_user_by_id(session, user_id) if not db_user: @@ -217,6 +223,21 @@ async def start_command_handler(message: types.Message, f"Failed to update existing user {user_id} in session: {e_update}", exc_info=True) + # Attribute user to ad campaign if start param provided + if ad_start_param: + try: + from db.dal import ad_dal as _ad_dal + campaign = await _ad_dal.get_campaign_by_start_param(session, ad_start_param) + if campaign and campaign.is_active: + await _ad_dal.ensure_attribution(session, user_id=user_id, campaign_id=campaign.ad_campaign_id) + await session.commit() + except Exception as e_attr: + logging.error(f"Failed to attribute user {user_id} to ad '{ad_start_param}': {e_attr}") + try: + await session.rollback() + except Exception: + pass + # Send welcome message if not disabled if not settings.DISABLE_WELCOME_MESSAGE: await message.answer(_(key="welcome", user_name=hd.quote(user.full_name))) diff --git a/bot/handlers/user/trial_handler.py b/bot/handlers/user/trial_handler.py index b03966c..c3ad911 100644 --- a/bot/handlers/user/trial_handler.py +++ b/bot/handlers/user/trial_handler.py @@ -113,6 +113,14 @@ async def request_trial_confirmation_handler( # Send notification to admin about new trial notification_service = NotificationService(callback.bot, settings, i18n) await notification_service.notify_trial_activation(user_id, end_date_obj) + # Mark ad attribution trial if exists + try: + from db.dal import ad_dal as _ad_dal + await _ad_dal.mark_trial_activated(session, user_id) + await session.commit() + except Exception as e_mark: + await session.rollback() + logging.error(f"Failed to mark trial for ad attribution for user {user_id}: {e_mark}") else: message_key_from_service = ( activation_result.get("message_key", "trial_activation_failed") @@ -298,6 +306,13 @@ async def confirm_activate_trial_handler( if activation_result and activation_result.get("activated") and end_date_obj: notification_service = NotificationService(callback.bot, settings, i18n) await notification_service.notify_trial_activation(user_id, end_date_obj) + try: + from db.dal import ad_dal as _ad_dal + await _ad_dal.mark_trial_activated(session, user_id) + await session.commit() + except Exception as e_mark: + await session.rollback() + logging.error(f"Failed to mark trial for ad attribution for user {user_id}: {e_mark}") @router.callback_query(F.data == "main_action:cancel_trial") diff --git a/bot/keyboards/inline/admin_keyboards.py b/bot/keyboards/inline/admin_keyboards.py index 3fa2ccc..85677dd 100644 --- a/bot/keyboards/inline/admin_keyboards.py +++ b/bot/keyboards/inline/admin_keyboards.py @@ -25,6 +25,10 @@ def get_admin_panel_keyboard(i18n_instance, lang: str, builder.button(text=_(key="admin_promo_marketing_section"), callback_data="admin_section:promo_marketing") + # Реклама + builder.button(text=_(key="admin_ads_section", default="📈 Реклама"), + callback_data="admin_action:ads") + # Системные функции builder.button(text=_(key="admin_system_functions_section"), callback_data="admin_section:system_functions") @@ -116,6 +120,17 @@ def get_system_functions_keyboard(i18n_instance, lang: str) -> InlineKeyboardMar return builder.as_markup() +def get_ads_menu_keyboard(i18n_instance, lang: str) -> InlineKeyboardMarkup: + _ = lambda key, **kwargs: i18n_instance.gettext(lang, key, **kwargs) + builder = InlineKeyboardBuilder() + builder.button(text=_(key="admin_ads_create_button", default="➕ Создать кампанию"), + callback_data="admin_action:ads_create") + builder.button(text=_(key="back_to_admin_panel_button"), + callback_data="admin_action:main") + builder.adjust(1, 1) + return builder.as_markup() + + def get_logs_menu_keyboard(i18n_instance, lang: str) -> InlineKeyboardMarkup: _ = lambda key, **kwargs: i18n_instance.gettext(lang, key, **kwargs) builder = InlineKeyboardBuilder() diff --git a/bot/states/admin_states.py b/bot/states/admin_states.py index 330e53e..6db5421 100644 --- a/bot/states/admin_states.py +++ b/bot/states/admin_states.py @@ -28,3 +28,8 @@ class AdminStates(StatesGroup): waiting_for_user_search = State() waiting_for_subscription_days_to_add = State() waiting_for_direct_message_to_user = State() + + # Ads campaigns + waiting_for_ad_source = State() + waiting_for_ad_start_param = State() + waiting_for_ad_cost = State() diff --git a/db/dal/__init__.py b/db/dal/__init__.py new file mode 100644 index 0000000..8df65b5 --- /dev/null +++ b/db/dal/__init__.py @@ -0,0 +1,21 @@ +from . import user_dal +from . import payment_dal +from . import subscription_dal +from . import promo_code_dal +from . import panel_sync_dal +from . import message_log_dal +from . import user_billing_dal +from . import ad_dal + +__all__ = ( + "user_dal", + "payment_dal", + "subscription_dal", + "promo_code_dal", + "panel_sync_dal", + "message_log_dal", + "user_billing_dal", + "ad_dal", +) + + diff --git a/db/dal/ad_dal.py b/db/dal/ad_dal.py new file mode 100644 index 0000000..85f820c --- /dev/null +++ b/db/dal/ad_dal.py @@ -0,0 +1,129 @@ +import logging +from typing import Optional, List, Dict, Any, Tuple +from sqlalchemy.ext.asyncio import AsyncSession +from sqlalchemy.future import select +from sqlalchemy.orm import selectinload +from sqlalchemy import update, delete, func, and_ + +from ..models import AdCampaign, AdAttribution, Payment + + +async def create_campaign( + session: AsyncSession, *, source: str, start_param: str, cost: float +) -> AdCampaign: + existing = await get_campaign_by_start_param(session, start_param) + if existing: + raise ValueError("ad_campaign_start_param_exists") + + campaign = AdCampaign(source=source, start_param=start_param, cost=float(cost)) + session.add(campaign) + await session.flush() + await session.refresh(campaign) + logging.info( + f"AdCampaign created id={campaign.ad_campaign_id}, source={source}, start={start_param}, cost={cost}" + ) + return campaign + + +async def get_campaign_by_id(session: AsyncSession, campaign_id: int) -> Optional[AdCampaign]: + stmt = select(AdCampaign).where(AdCampaign.ad_campaign_id == campaign_id) + result = await session.execute(stmt) + return result.scalar_one_or_none() + + +async def get_campaign_by_start_param(session: AsyncSession, start_param: str) -> Optional[AdCampaign]: + clean = start_param.strip() + stmt = select(AdCampaign).where(AdCampaign.start_param == clean) + result = await session.execute(stmt) + return result.scalar_one_or_none() + + +async def list_campaigns(session: AsyncSession, *, only_active: bool = False) -> List[AdCampaign]: + stmt = select(AdCampaign).order_by(AdCampaign.created_at.desc()) + if only_active: + stmt = stmt.where(AdCampaign.is_active == True) + result = await session.execute(stmt) + return result.scalars().all() + + +async def toggle_campaign_active(session: AsyncSession, campaign_id: int, is_active: bool) -> bool: + stmt = ( + update(AdCampaign) + .where(AdCampaign.ad_campaign_id == campaign_id) + .values(is_active=is_active) + ) + result = await session.execute(stmt) + return result.rowcount > 0 + + +async def ensure_attribution(session: AsyncSession, *, user_id: int, campaign_id: int) -> AdAttribution: + existing = await get_attribution_for_user(session, user_id) + if existing: + return existing + attrib = AdAttribution(user_id=user_id, ad_campaign_id=campaign_id) + session.add(attrib) + await session.flush() + await session.refresh(attrib) + logging.info(f"AdAttribution created for user {user_id} -> campaign {campaign_id}") + return attrib + + +async def get_attribution_for_user(session: AsyncSession, user_id: int) -> Optional[AdAttribution]: + stmt = select(AdAttribution).where(AdAttribution.user_id == user_id) + result = await session.execute(stmt) + return result.scalar_one_or_none() + + +async def mark_trial_activated(session: AsyncSession, user_id: int) -> bool: + stmt = ( + update(AdAttribution) + .where(and_(AdAttribution.user_id == user_id, AdAttribution.trial_activated_at.is_(None))) + .values(trial_activated_at=func.now()) + ) + result = await session.execute(stmt) + return result.rowcount > 0 + + +async def get_campaign_stats(session: AsyncSession, campaign_id: int) -> Dict[str, Any]: + # Starts (attributed users) + starts_stmt = select(func.count(AdAttribution.user_id)).where( + AdAttribution.ad_campaign_id == campaign_id + ) + starts = (await session.execute(starts_stmt)).scalar() or 0 + + # Trials + trials_stmt = select(func.count(AdAttribution.user_id)).where( + and_(AdAttribution.ad_campaign_id == campaign_id, AdAttribution.trial_activated_at.is_not(None)) + ) + trials = (await session.execute(trials_stmt)).scalar() or 0 + + # Payers (unique users with succeeded payments) + payers_stmt = select(func.count(func.distinct(Payment.user_id))).select_from(Payment).where( + and_( + Payment.status == "succeeded", + Payment.user_id.in_( + select(AdAttribution.user_id).where(AdAttribution.ad_campaign_id == campaign_id) + ), + ) + ) + payers = (await session.execute(payers_stmt)).scalar() or 0 + + # Revenue sum + revenue_stmt = select(func.coalesce(func.sum(Payment.amount), 0.0)).select_from(Payment).where( + and_( + Payment.status == "succeeded", + Payment.user_id.in_( + select(AdAttribution.user_id).where(AdAttribution.ad_campaign_id == campaign_id) + ), + ) + ) + revenue = float((await session.execute(revenue_stmt)).scalar() or 0.0) + + return { + "starts": int(starts), + "trials": int(trials), + "payers": int(payers), + "revenue": revenue, + } + + diff --git a/db/models.py b/db/models.py index 8f661e6..d33dfdf 100644 --- a/db/models.py +++ b/db/models.py @@ -227,3 +227,35 @@ class PanelSyncStatus(Base): subscriptions_synced = Column(Integer, default=0) __table_args__ = (UniqueConstraint('id'), ) + + +class AdCampaign(Base): + __tablename__ = "ad_campaigns" + + ad_campaign_id = Column(Integer, primary_key=True, autoincrement=True) + source = Column(String, nullable=False, index=True) + start_param = Column(String, nullable=False, unique=True, index=True) + cost = Column(Float, nullable=False, default=0.0) + is_active = Column(Boolean, default=True, index=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + attributions = relationship( + "AdAttribution", + back_populates="campaign", + cascade="all, delete-orphan", + ) + + def __repr__(self): + return f"" + + +class AdAttribution(Base): + __tablename__ = "ad_attributions" + + user_id = Column(BigInteger, ForeignKey("users.user_id"), primary_key=True, index=True) + ad_campaign_id = Column(Integer, ForeignKey("ad_campaigns.ad_campaign_id"), nullable=False, index=True) + first_start_at = Column(DateTime(timezone=True), server_default=func.now()) + trial_activated_at = Column(DateTime(timezone=True), nullable=True) + + user = relationship("User") + campaign = relationship("AdCampaign", back_populates="attributions") diff --git a/locales/en.json b/locales/en.json index fd11ac5..9275745 100644 --- a/locales/en.json +++ b/locales/en.json @@ -406,5 +406,19 @@ "error_creating_payment_record": "Error creating payment record. Please try again later.", "error_payment_gateway_link_failed": "Error creating payment link. Please try again later.", "status_active": "Active", - "status_inactive": "Inactive" + "status_inactive": "Inactive", + "admin_ads_section": "📈 Ads", + "admin_ads_header": "📈 Ad Campaigns:", + "admin_ads_empty": "📭 No ad campaigns. Click \"Create\" to add one.", + "admin_ads_item": "ID: {id}\nSource: {source}\nstart={start_param}\nCost: {cost} RUB\nActive: {active}\n— Starts: {starts}\n— Trials: {trials}\n— Payers: {payers}\n— Revenue: {revenue} RUB", + "admin_ads_create_button": "➕ Create campaign", + "admin_ads_create_source_prompt": "Enter source (e.g., AEZA, VK, TG-channel):", + "admin_ads_create_start_param_prompt": "Enter start link parameter (e.g., AEZA). Will be used as start=AEZA", + "admin_ads_create_cost_prompt": "Enter campaign cost (RUB):", + "admin_ads_invalid_source": "❌ Invalid source. Enter up to 64 characters.", + "admin_ads_invalid_start_param": "❌ Invalid parameter. Allowed letters/digits/underscore/dash (2-64).", + "admin_ads_invalid_cost": "❌ Invalid amount. Enter a non-negative number.", + "admin_ads_start_param_exists": "❌ A campaign with this start parameter already exists.", + "admin_ads_created_success": "✅ Campaign created!\nID: {id}\nSource: {source}\nParam: {start_param}\nCost: {cost} RUB", + "admin_ads_back_to_menu_hint": "Done. Back to Ads section:" } diff --git a/locales/ru.json b/locales/ru.json index efa4954..990fdea 100644 --- a/locales/ru.json +++ b/locales/ru.json @@ -405,5 +405,19 @@ "error_creating_payment_record": "Ошибка создания записи платежа. Попробуйте позже.", "error_payment_gateway_link_failed": "Ошибка создания платежной ссылки. Попробуйте позже.", "status_active": "Активна", - "status_inactive": "Неактивна" + "status_inactive": "Неактивна", + "admin_ads_section": "📈 Реклама", + "admin_ads_header": "📈 Рекламные кампании:", + "admin_ads_empty": "📭 Рекламные кампании отсутствуют. Нажмите \"Создать\" чтобы добавить новую.", + "admin_ads_item": "ID: {id}\nИсточник: {source}\nstart={start_param}\nСтоимость: {cost} RUB\nАктивна: {active}\n— Запустили: {starts}\n— Взяли триал: {trials}\n— Оплатили: {payers}\n— Доход: {revenue} RUB", + "admin_ads_create_button": "➕ Создать кампанию", + "admin_ads_create_source_prompt": "Введите источник (например: AEZA, VK, TG-канал):", + "admin_ads_create_start_param_prompt": "Введите параметр старт-ссылки (например: AEZA). Будет использован как start=AEZA", + "admin_ads_create_cost_prompt": "Введите сумму затрат на кампанию (в RUB):", + "admin_ads_invalid_source": "❌ Неверный источник. Введите до 64 символов.", + "admin_ads_invalid_start_param": "❌ Неверный параметр. Допустимы буквы/цифры/подчёркивания/дефисы (2-64).", + "admin_ads_invalid_cost": "❌ Неверная сумма. Введите неотрицательное число.", + "admin_ads_start_param_exists": "❌ Кампания с таким start-параметром уже существует.", + "admin_ads_created_success": "✅ Кампания создана!\nID: {id}\nИсточник: {source}\nПараметр: {start_param}\nЗатраты: {cost} RUB", + "admin_ads_back_to_menu_hint": "Готово. Вернуться к разделу рекламы:" }