From 443e2e62db58ff1508daf274c653cbf60a181eda Mon Sep 17 00:00:00 2001 From: 3252a8 <3252a8@proton.me> Date: Fri, 22 May 2026 22:39:45 +0300 Subject: [PATCH] fix: auto-merge duplicate panel identities --- backend/bot/handlers/admin/sync_admin.py | 278 ++++++++++++++++++++++- tests/test_admin_sync_performance.py | 98 ++++++++ 2 files changed, 374 insertions(+), 2 deletions(-) diff --git a/backend/bot/handlers/admin/sync_admin.py b/backend/bot/handlers/admin/sync_admin.py index c56b507..453d206 100644 --- a/backend/bot/handlers/admin/sync_admin.py +++ b/backend/bot/handlers/admin/sync_admin.py @@ -1,6 +1,6 @@ import asyncio import logging -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from typing import Any, Optional, Union from aiogram import Bot, Router, types @@ -108,6 +108,21 @@ def _as_utc(value: datetime) -> datetime: return value.astimezone(timezone.utc) +def _panel_expire_at(panel_user: dict[str, Any]) -> Optional[datetime]: + raw_value = panel_user.get("expireAt") + if not raw_value: + return None + try: + return datetime.fromisoformat(str(raw_value).replace("Z", "+00:00")) + except (TypeError, ValueError): + return None + + +def _panel_subscription_uuid(panel_user: dict[str, Any]) -> Optional[str]: + value = panel_user.get("subscriptionUuid") or panel_user.get("shortUuid") + return str(value) if value else None + + def _should_update_lifetime_used_traffic( existing_user, lifetime_used: int, @@ -350,6 +365,211 @@ async def _bind_panel_email_to_user( return existing_user, True +async def _merge_local_duplicate_panel_user_if_needed( + session: AsyncSession, + *, + existing_user, + duplicate_panel_uuid: str, +): + duplicate_local_user = await user_dal.get_user_by_panel_uuid(session, duplicate_panel_uuid) + if not duplicate_local_user or duplicate_local_user.user_id == existing_user.user_id: + return existing_user, True + + try: + merged_user = await user_dal.merge_users( + session, + source_user_id=duplicate_local_user.user_id, + target_user_id=existing_user.user_id, + ) + logging.info( + "Sync: merged local duplicate user %s into %s for duplicate panel UUID %s.", + duplicate_local_user.user_id, + merged_user.user_id, + duplicate_panel_uuid, + ) + return merged_user, True + except Exception as exc: + logging.warning( + "Sync: could not merge local duplicate user %s into %s for panel UUID %s: %s", + duplicate_local_user.user_id, + existing_user.user_id, + duplicate_panel_uuid, + exc, + ) + return existing_user, False + + +def _panel_identity_payload_with_expiry( + user, + *, + expire_at: datetime, +) -> dict[str, Any]: + description_text = "\n".join( + line + for line in [ + user.email or "", + user.username or "", + user.first_name or "", + user.last_name or "", + ] + if line + ) + payload = _panel_identity_update_payload(user, description_text) + payload["expireAt"] = expire_at.isoformat(timespec="milliseconds").replace("+00:00", "Z") + if expire_at > datetime.now(timezone.utc): + payload["status"] = "ACTIVE" + return payload + + +async def _absorb_duplicate_panel_identity( + session: AsyncSession, + *, + panel_service: PanelApiService, + existing_user, + keep_panel_uuid: str, + keep_panel_user: Optional[dict[str, Any]], + duplicate_panel_user: dict[str, Any], + settings: Settings, + subscriptions_by_panel_uuid: dict[str, Subscription], + active_subscriptions_by_user_panel: dict[tuple[int, str], Subscription], +) -> dict[str, int | bool]: + duplicate_panel_uuid = str(duplicate_panel_user.get("uuid") or "") + if not duplicate_panel_uuid: + return {"resolved": False, "subscriptions_created": 0, "subscriptions_updated": 0} + + subscriptions_created = 0 + subscriptions_updated = 0 + now = datetime.now(timezone.utc) + duplicate_expire_at = _panel_expire_at(duplicate_panel_user) + duplicate_status = str(duplicate_panel_user.get("status") or "").upper() + duplicate_is_active = bool( + duplicate_expire_at and duplicate_status == "ACTIVE" and duplicate_expire_at > now + ) + + keep_subscription_uuid = _panel_subscription_uuid(keep_panel_user or {}) + target_sub = ( + subscriptions_by_panel_uuid.get(keep_subscription_uuid) if keep_subscription_uuid else None + ) + if not target_sub: + target_sub = active_subscriptions_by_user_panel.get( + (int(existing_user.user_id), keep_panel_uuid) + ) + + final_end_date: Optional[datetime] = None + if duplicate_is_active and duplicate_expire_at: + source_remaining = max(timedelta(0), duplicate_expire_at - now) + if target_sub: + target_end = _as_utc(target_sub.end_date) + base_end = target_end if target_end > now else now + final_end_date = base_end + source_remaining + update_payload: dict[str, Any] = { + "user_id": int(existing_user.user_id), + "panel_user_uuid": keep_panel_uuid, + "end_date": final_end_date, + "is_active": True, + "status_from_panel": "ACTIVE_EXTENDED_BY_PANEL_DUPLICATE_MERGE", + } + if keep_subscription_uuid: + update_payload["panel_subscription_uuid"] = keep_subscription_uuid + update_delta = _subscription_update_delta(target_sub, update_payload) + if update_delta: + await subscription_dal.update_subscription( + session, + target_sub.subscription_id, + update_delta, + ) + for key, value in update_delta.items(): + setattr(target_sub, key, value) + subscriptions_updated += 1 + elif keep_subscription_uuid: + final_end_date = now + (duplicate_expire_at - now) + created_sub = await subscription_dal.upsert_subscription( + session, + { + "user_id": int(existing_user.user_id), + "panel_user_uuid": keep_panel_uuid, + "panel_subscription_uuid": keep_subscription_uuid, + "start_date": None, + "end_date": final_end_date, + "duration_months": None, + "is_active": True, + "status_from_panel": "ACTIVE_EXTENDED_BY_PANEL_DUPLICATE_MERGE", + "traffic_limit_bytes": getattr(settings, "user_traffic_limit_bytes", 0), + "auto_renew_enabled": False, + }, + ) + subscriptions_by_panel_uuid[keep_subscription_uuid] = created_sub + active_subscriptions_by_user_panel[ + (int(created_sub.user_id), created_sub.panel_user_uuid) + ] = created_sub + subscriptions_created += 1 + + duplicate_subscription_uuid = _panel_subscription_uuid(duplicate_panel_user) + duplicate_sub = ( + subscriptions_by_panel_uuid.get(duplicate_subscription_uuid) + if duplicate_subscription_uuid + else None + ) + if duplicate_sub and duplicate_sub is not target_sub: + await subscription_dal.update_subscription( + session, + duplicate_sub.subscription_id, + { + "user_id": int(existing_user.user_id), + "is_active": False, + "skip_notifications": True, + "status_from_panel": "MERGED_PANEL_DUPLICATE", + }, + ) + duplicate_sub.user_id = int(existing_user.user_id) + duplicate_sub.is_active = False + duplicate_sub.skip_notifications = True + duplicate_sub.status_from_panel = "MERGED_PANEL_DUPLICATE" + subscriptions_updated += 1 + elif not duplicate_sub: + await session.execute( + update(Subscription) + .where(Subscription.panel_user_uuid == duplicate_panel_uuid) + .values( + user_id=int(existing_user.user_id), + is_active=False, + skip_notifications=True, + status_from_panel="MERGED_PANEL_DUPLICATE", + ) + ) + + if final_end_date: + await panel_service.update_user_details_on_panel( + keep_panel_uuid, + _panel_identity_payload_with_expiry(existing_user, expire_at=final_end_date), + log_response=False, + ) + + deleted = await panel_service.delete_user_from_panel( + duplicate_panel_uuid, + log_response=False, + ) + if deleted: + logging.info( + "Sync: absorbed duplicate panel UUID %s into kept panel UUID %s for user %s.", + duplicate_panel_uuid, + keep_panel_uuid, + existing_user.user_id, + ) + else: + logging.warning( + "Sync: failed to delete duplicate panel UUID %s after absorbing it into %s.", + duplicate_panel_uuid, + keep_panel_uuid, + ) + + return { + "resolved": bool(deleted), + "subscriptions_created": subscriptions_created, + "subscriptions_updated": subscriptions_updated, + } + + async def perform_sync( panel_service: PanelApiService, session: AsyncSession, @@ -430,6 +650,11 @@ async def _perform_sync_impl( subscriptions_by_panel_uuid = sync_indexes["subscriptions_by_panel_uuid"] active_subscriptions_by_user_panel = sync_indexes["active_subscriptions_by_user_panel"] panel_uuids_by_telegram_id = sync_indexes["panel_uuids_by_telegram_id"] + panel_users_by_uuid = { + str(panel_user["uuid"]): panel_user + for panel_user in panel_users_data + if panel_user.get("uuid") + } for panel_user_dict in panel_users_data: try: @@ -579,12 +804,61 @@ async def _perform_sync_impl( ) if linked_uuid_still_present: is_duplicate_panel_identity = True + ( + existing_user, + can_absorb_duplicate_panel_user, + ) = await _merge_local_duplicate_panel_user_if_needed( + session, + existing_user=existing_user, + duplicate_panel_uuid=panel_uuid, + ) + if not can_absorb_duplicate_panel_user: + logging.warning( + "Sync: duplicate panel users share telegramId %s; keeping local panel UUID %s and skipping duplicate panel UUID %s because local duplicate merge failed.", # noqa: E501 + telegram_id_from_panel, + linked_uuid, + panel_uuid, + ) + continue + actual_user_id = existing_user.user_id + users_by_panel_uuid[linked_uuid] = existing_user + if existing_user.telegram_id is not None: + users_by_telegram_id[int(existing_user.telegram_id)] = existing_user + users_by_user_id[int(existing_user.user_id)] = existing_user + if existing_user.email: + users_by_email[existing_user.email.strip().lower()] = existing_user + merge_result = await _absorb_duplicate_panel_identity( + session, + panel_service=panel_service, + existing_user=existing_user, + keep_panel_uuid=str(linked_uuid), + keep_panel_user=panel_users_by_uuid.get(str(linked_uuid)), + duplicate_panel_user=panel_user_dict, + settings=settings, + subscriptions_by_panel_uuid=subscriptions_by_panel_uuid, + active_subscriptions_by_user_panel=( + active_subscriptions_by_user_panel + ), + ) + subscriptions_created += int(merge_result["subscriptions_created"]) + subscriptions_updated += int(merge_result["subscriptions_updated"]) + subscriptions_synced_count += int( + merge_result["subscriptions_created"] + ) + int(merge_result["subscriptions_updated"]) + if merge_result["resolved"]: + users_updated += 1 + users_uuid_updated += 1 + panel_uuids_by_telegram_id.get(telegram_id_from_panel, set()).discard( + str(panel_uuid) + ) + users_by_panel_uuid.pop(str(panel_uuid), None) logging.warning( - "Sync: duplicate panel users share telegramId %s; keeping local panel UUID %s and skipping duplicate panel UUID %s.", # noqa: E501 + "Sync: duplicate panel users share telegramId %s; kept local panel UUID %s and processed duplicate panel UUID %s.", # noqa: E501 telegram_id_from_panel, linked_uuid, panel_uuid, ) + continue else: existing_user.panel_user_uuid = panel_uuid user_was_updated = True diff --git a/tests/test_admin_sync_performance.py b/tests/test_admin_sync_performance.py index 94183ca..ce9e156 100644 --- a/tests/test_admin_sync_performance.py +++ b/tests/test_admin_sync_performance.py @@ -1,7 +1,10 @@ +import asyncio from datetime import datetime, timedelta, timezone from types import SimpleNamespace +from unittest.mock import AsyncMock, patch from bot.handlers.admin.sync_admin import ( + _absorb_duplicate_panel_identity, _coerce_panel_telegram_id, _description_matches, _should_update_lifetime_used_traffic, @@ -136,3 +139,98 @@ def test_lifetime_traffic_update_allows_large_delta_and_skips_duplicate_panel_id settings=settings, is_duplicate_panel_identity=True, ) + + +def test_absorb_duplicate_panel_identity_extends_kept_user_and_deletes_duplicate(): + now = datetime.now(timezone.utc) + target_sub = SimpleNamespace( + subscription_id=10, + user_id=42, + panel_user_uuid="panel-keep", + panel_subscription_uuid="sub-keep", + end_date=now - timedelta(days=2), + is_active=False, + status_from_panel="EXPIRED", + ) + duplicate_sub = SimpleNamespace( + subscription_id=11, + user_id=42, + panel_user_uuid="panel-duplicate", + panel_subscription_uuid="sub-duplicate", + end_date=now + timedelta(days=30), + is_active=True, + skip_notifications=False, + status_from_panel="ACTIVE", + ) + panel_service = SimpleNamespace( + update_user_details_on_panel=AsyncMock(return_value={"uuid": "panel-keep"}), + delete_user_from_panel=AsyncMock(return_value=True), + ) + session = SimpleNamespace(execute=AsyncMock()) + settings = SimpleNamespace(user_traffic_limit_bytes=0) + user = SimpleNamespace( + user_id=42, + panel_user_uuid="panel-keep", + telegram_id=969808056, + email="paid@example.com", + username="alice", + first_name="Alice", + last_name=None, + ) + + async def update_subscription(_session, subscription_id, update_data): + sub = target_sub if subscription_id == target_sub.subscription_id else duplicate_sub + for key, value in update_data.items(): + setattr(sub, key, value) + return sub + + with patch( + "bot.handlers.admin.sync_admin.subscription_dal.update_subscription", + AsyncMock(side_effect=update_subscription), + ): + result = asyncio.run( + _absorb_duplicate_panel_identity( + session, + panel_service=panel_service, + existing_user=user, + keep_panel_uuid="panel-keep", + keep_panel_user={ + "uuid": "panel-keep", + "subscriptionUuid": "sub-keep", + "status": "EXPIRED", + "expireAt": (now - timedelta(days=2)).isoformat(), + }, + duplicate_panel_user={ + "uuid": "panel-duplicate", + "subscriptionUuid": "sub-duplicate", + "telegramId": 969808056, + "status": "ACTIVE", + "expireAt": (now + timedelta(days=30)).isoformat(), + }, + settings=settings, + subscriptions_by_panel_uuid={ + "sub-keep": target_sub, + "sub-duplicate": duplicate_sub, + }, + active_subscriptions_by_user_panel={}, + ) + ) + + assert result["resolved"] + assert result["subscriptions_updated"] == 2 + assert target_sub.is_active + assert target_sub.status_from_panel == "ACTIVE_EXTENDED_BY_PANEL_DUPLICATE_MERGE" + assert target_sub.panel_user_uuid == "panel-keep" + assert target_sub.end_date > now + timedelta(days=29) + assert not duplicate_sub.is_active + assert duplicate_sub.skip_notifications + assert duplicate_sub.status_from_panel == "MERGED_PANEL_DUPLICATE" + panel_service.update_user_details_on_panel.assert_awaited_once() + update_uuid, update_payload = panel_service.update_user_details_on_panel.await_args.args[:2] + assert update_uuid == "panel-keep" + assert update_payload["status"] == "ACTIVE" + assert update_payload["telegramId"] == 969808056 + panel_service.delete_user_from_panel.assert_awaited_once_with( + "panel-duplicate", + log_response=False, + )