From 99648a956e20cb99d406455db2f26db469ba7235 Mon Sep 17 00:00:00 2001 From: Fringg Date: Mon, 4 May 2026 05:27:49 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20persist=20campaign=20across=20bot?= =?UTF-8?q?=E2=86=92webapp=20registration=20handoff=20via=20Redis?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/cabinet/routes/auth.py | 162 +++++++++++++++++++------------ app/handlers/start.py | 48 ++++++++- app/services/referral_service.py | 67 +++++++++++++ 3 files changed, 216 insertions(+), 61 deletions(-) diff --git a/app/cabinet/routes/auth.py b/app/cabinet/routes/auth.py index a590b6da..7f30cba8 100644 --- a/app/cabinet/routes/auth.py +++ b/app/cabinet/routes/auth.py @@ -172,74 +172,112 @@ async def _process_campaign_bonus( db: AsyncSession, user: User, campaign_slug: str | None, + telegram_id: int | None = None, ) -> CampaignBonusInfo | None: - """Process campaign bonus for user during auth. Never raises.""" + """Process campaign bonus for user during auth. Never raises. + + If ``campaign_slug`` is not provided but ``telegram_id`` is given, the + function falls back to Redis ``pending_campaign:{telegram_id}`` -- populated + by the bot's /start handler when a user opens an advertising campaign link + but then completes registration via the cabinet WebApp (Telegram menu + button) instead of the bot dialog. The Redis entry is cleared after a + successful consumption attempt. + """ + pending_campaign_consumed = False + if not campaign_slug and telegram_id: + try: + from app.services.referral_service import get_pending_campaign + + pending = await get_pending_campaign(telegram_id) + if pending and pending.get('campaign_slug'): + campaign_slug = pending['campaign_slug'] + pending_campaign_consumed = True + logger.info( + 'Resolved campaign from Redis pending_campaign (cabinet)', + telegram_id=telegram_id, + campaign_slug=campaign_slug, + ) + except Exception as e: + logger.warning('Failed to check pending campaign', error=e) + 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 + try: + campaign = await get_campaign_by_start_parameter(db, campaign_slug, only_active=True) + if not campaign: + return None - # Skip if user IS the campaign partner — prevent self-referral - if campaign.partner_user_id and campaign.partner_user_id == user.id: - logger.debug( - 'Skipping campaign attribution: user is the campaign partner', - user_id=user.id, - campaign_id=campaign.id, - ) - 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 - - # Привязать реферала к партнёру кампании (если партнёр назначен и юзер ещё не привязан) - if campaign.partner_user_id and not user.referred_by_id: - user.referred_by_id = campaign.partner_user_id - await db.flush() - try: - from app.bot_factory import create_bot - - async with create_bot() as bot: - await process_referral_registration(db, user.id, campaign.partner_user_id, bot=bot) - logger.info( - 'Referral set from campaign partner', + # Skip if user IS the campaign partner — prevent self-referral + if campaign.partner_user_id and campaign.partner_user_id == user.id: + logger.debug( + 'Skipping campaign attribution: user is the campaign partner', user_id=user.id, - partner_user_id=campaign.partner_user_id, campaign_id=campaign.id, ) - except Exception as e: - logger.error('Failed to process referral from campaign partner', error=e) + return None - service = AdvertisingCampaignService() - result = await service.apply_campaign_bonus(db, user, campaign) - if not result.success: - 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()) - # Refresh user to get updated balance after bonus - await db.refresh(user) + 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 - 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 + # Привязать реферала к партнёру кампании (если партнёр назначен и юзер ещё не привязан) + if campaign.partner_user_id and not user.referred_by_id: + user.referred_by_id = campaign.partner_user_id + await db.flush() + try: + from app.bot_factory import create_bot + + async with create_bot() as bot: + await process_referral_registration(db, user.id, campaign.partner_user_id, bot=bot) + logger.info( + 'Referral set from campaign partner', + user_id=user.id, + partner_user_id=campaign.partner_user_id, + campaign_id=campaign.id, + ) + except Exception as e: + logger.error('Failed to process referral from campaign partner', error=e) + + 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 rollback after campaign bonus error', user_id=user.id) - return None + 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 + finally: + # Clear Redis pending_campaign whenever we consumed it. Done regardless + # of success — if processing failed (already applied, race, exception), + # we don't want to keep retrying on every subsequent login. + if pending_campaign_consumed and telegram_id: + try: + from app.services.referral_service import clear_pending_campaign + + await clear_pending_campaign(telegram_id) + except Exception: + pass async def _process_referral_code( @@ -584,8 +622,11 @@ async def auth_telegram( except Exception: pass - # Process campaign bonus - response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug) + # Process campaign bonus. + # Pass telegram_id so the function can fall back to Redis pending_campaign + # if the user came via /start in the bot but completed + # registration in the WebApp without an explicit campaign_slug. + response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug, telegram_id=telegram_id) if response.campaign_bonus: response.user = _user_to_response(user) @@ -691,8 +732,8 @@ async def auth_telegram_widget( except Exception: pass - # Process campaign bonus - response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug) + # Process campaign bonus (pending_campaign Redis fallback for Telegram Login Widget) + response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug, telegram_id=request.id) if response.campaign_bonus: response.user = _user_to_response(user) @@ -837,7 +878,8 @@ async def auth_telegram_oidc( except Exception: pass - response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug) + # Process campaign bonus (pending_campaign Redis fallback for Telegram OIDC) + response.campaign_bonus = await _process_campaign_bonus(db, user, request.campaign_slug, telegram_id=telegram_id) if response.campaign_bonus: response.user = _user_to_response(user) diff --git a/app/handlers/start.py b/app/handlers/start.py index da227091..1ab4b3ed 100644 --- a/app/handlers/start.py +++ b/app/handlers/start.py @@ -49,7 +49,11 @@ from app.services.pinned_message_service import ( get_active_pinned_message, ) from app.services.privacy_policy_service import PrivacyPolicyService -from app.services.referral_service import process_referral_registration, save_pending_referral +from app.services.referral_service import ( + process_referral_registration, + save_pending_campaign, + save_pending_referral, +) from app.services.subscription_service import SubscriptionService from app.services.support_settings_service import SupportSettingsService from app.services.web_auth_service import WEB_AUTH_TOKEN_MIN_LENGTH, link_web_auth_token @@ -386,6 +390,17 @@ async def _apply_campaign_bonus_if_needed( if not result.success: return None + # Bot-flow successfully applied the campaign — clear the Redis pending entry + # (set in cmd_start as a fallback for the cabinet WebApp path) so it isn't + # re-evaluated on a subsequent cabinet login. + try: + from app.services.referral_service import clear_pending_campaign + + if getattr(user, 'telegram_id', None): + await clear_pending_campaign(user.telegram_id) + except Exception: + pass + if result.bonus_type == 'balance': amount_text = texts.format_price(result.balance_kopeks) return texts.CAMPAIGN_BONUS_BALANCE.format( @@ -724,6 +739,23 @@ async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession, start_parameter=campaign.start_parameter, ) await state.update_data(campaign_id=campaign.id) + # Persist campaign to Redis immediately so it survives if user opens + # miniapp/cabinet (via Telegram menu button) before completing the + # bot registration flow. Mirrors the pending_referral mechanism. + # Only for new users — existing users already had attribution applied. + if not db_user: + try: + await save_pending_campaign( + message.from_user.id, + campaign.start_parameter, + campaign.id, + ) + except Exception as exc: + logger.warning( + 'Failed to persist pending campaign', + campaign_id=campaign.id, + error=exc, + ) if campaign.partner_user_id: await state.update_data(referrer_id=campaign.partner_user_id) logger.info( @@ -2416,6 +2448,20 @@ async def required_sub_channel_check( campaign_id=campaign.id, partner_user_id=campaign.partner_user_id, ) + # Mirror save in Redis so cabinet WebApp auth can pick it up + # if user opens miniapp before completing registration. + try: + await save_pending_campaign( + query.from_user.id, + campaign.start_parameter, + campaign.id, + ) + except Exception as exc: + logger.warning( + 'Failed to persist pending campaign after channel check', + campaign_id=campaign.id, + error=exc, + ) else: state_data['referral_code'] = pending_start_payload logger.info( diff --git a/app/services/referral_service.py b/app/services/referral_service.py index d08e6e4d..7409cce3 100644 --- a/app/services/referral_service.py +++ b/app/services/referral_service.py @@ -99,6 +99,73 @@ async def clear_pending_referral(telegram_id: int) -> None: pass +# --------------------------------------------------------------------------- +# Pending campaign helpers (Redis) +# +# Mirrors pending_referral: lets us survive the case where /start +# stored campaign_id only in FSM, but the user opened the cabinet WebApp before +# completing bot registration. The cabinet auth route reads this as a fallback +# when the HTTP request didn't carry an explicit campaign_slug. +# --------------------------------------------------------------------------- +_PENDING_CAMPAIGN_TTL = 7 * 24 * 3600 # 7 days + + +async def save_pending_campaign(telegram_id: int, campaign_slug: str, campaign_id: int) -> bool: + """Save pending campaign attribution to Redis for a not-yet-registered user. + + Called from /start handler immediately after resolving an advertising campaign. + Picked up by the cabinet auth route if the user opens the WebApp before + completing the bot registration flow. + """ + client = _get_redis() + if client is None: + return False + try: + key = f'pending_campaign:{telegram_id}' + data = json.dumps({'campaign_slug': campaign_slug, 'campaign_id': campaign_id}) + await client.setex(key, _PENDING_CAMPAIGN_TTL, data) + logger.info( + 'Saved pending campaign to Redis', + telegram_id=telegram_id, + campaign_slug=campaign_slug, + campaign_id=campaign_id, + ) + return True + except Exception as exc: + logger.warning('Failed to save pending campaign to Redis', error=exc) + return False + + +async def get_pending_campaign(telegram_id: int) -> dict[str, str | int] | None: + """Get pending campaign from Redis. + + Returns ``{'campaign_slug': ..., 'campaign_id': ...}`` or ``None``. + """ + client = _get_redis() + if client is None: + return None + try: + key = f'pending_campaign:{telegram_id}' + data = await client.get(key) + if data: + return json.loads(data) + return None + except Exception as exc: + logger.warning('Failed to get pending campaign from Redis', error=exc) + return None + + +async def clear_pending_campaign(telegram_id: int) -> None: + """Clear pending campaign after successful application.""" + client = _get_redis() + if client is None: + return + try: + await client.delete(f'pending_campaign:{telegram_id}') + except Exception: + pass + + async def _is_commission_limit_reached(db: AsyncSession, referrer_id: int, referral_id: int) -> bool: """Проверяет, исчерпан ли лимит комиссионных платежей для пары реферер-реферал.""" if settings.REFERRAL_MAX_COMMISSION_PAYMENTS <= 0: