Revert "Устранены дубли активных записей, ошибка уникальности и улучшена идемпотентность"

This commit is contained in:
machka pasla
2025-11-13 23:01:46 +03:00
committed by GitHub
parent 2b0fa3a314
commit cd816cbe5c
+67 -180
View File
@@ -3,7 +3,6 @@ from aiogram import Router, types, Bot
from aiogram.filters import Command
from typing import Optional, Union
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import update, or_
from datetime import datetime, timezone
from config.settings import Settings
@@ -11,19 +10,14 @@ from bot.services.panel_api_service import PanelApiService
from bot.services.notification_service import NotificationService
from db.dal import user_dal, subscription_dal, panel_sync_dal
from db.models import Subscription
from bot.middlewares.i18n import JsonI18n
router = Router(name="admin_sync_router")
async def perform_sync(
panel_service: PanelApiService,
session: AsyncSession,
settings: Settings,
i18n_instance: JsonI18n,
) -> dict:
async def perform_sync(panel_service: PanelApiService, session: AsyncSession,
settings: Settings, i18n_instance: JsonI18n) -> dict:
"""
Perform panel synchronization and return results
Returns dict with status, details, and sync statistics
@@ -33,7 +27,7 @@ async def perform_sync(
users_updated = 0
subscriptions_synced_count = 0
sync_errors = []
# Additional counters for detailed logging
users_without_telegram_id = 0
users_not_found_in_db = 0
@@ -58,12 +52,7 @@ async def perform_sync(
session, "success", status_msg, 0, 0
)
await session.commit()
return {
"status": "success",
"details": status_msg,
"users_synced": 0,
"subs_synced": 0,
}
return {"status": "success", "details": status_msg, "users_synced": 0, "subs_synced": 0}
total_panel_users = len(panel_users_data)
logging.info(f"Starting sync for {total_panel_users} panel users.")
@@ -72,16 +61,12 @@ async def perform_sync(
try:
panel_records_checked += 1
panel_uuid = panel_user_dict.get("uuid")
panel_subscription_uuid = panel_user_dict.get("subscriptionUuid") or panel_user_dict.get(
"shortUuid"
)
panel_subscription_uuid = panel_user_dict.get("subscriptionUuid") or panel_user_dict.get("shortUuid")
telegram_id_from_panel = panel_user_dict.get("telegramId")
if not panel_uuid:
sync_errors.append(f"Panel user missing UUID: {panel_user_dict}")
logging.warning(
f"Skipping panel user without UUID: {panel_user_dict}"
)
logging.warning(f"Skipping panel user without UUID: {panel_user_dict}")
continue
# Track users without telegram ID
@@ -90,35 +75,22 @@ async def perform_sync(
# Try to find existing user in local DB
existing_user = None
# First, try to find by telegram ID if available
if telegram_id_from_panel:
existing_user = await user_dal.get_user_by_id(
session, telegram_id_from_panel
)
existing_user = await user_dal.get_user_by_id(session, telegram_id_from_panel)
if existing_user:
logging.debug(
f"Found user by telegramId {telegram_id_from_panel}"
)
logging.debug(f"Found user by telegramId {telegram_id_from_panel}")
# If not found by telegram ID, try to find by panel UUID
if not existing_user:
existing_user = await user_dal.get_user_by_panel_uuid(
session, panel_uuid
)
existing_user = await user_dal.get_user_by_panel_uuid(session, panel_uuid)
if existing_user:
logging.info(
f"Found user by panel UUID {panel_uuid}, telegramId: {existing_user.user_id}"
)
logging.info(f"Found user by panel UUID {panel_uuid}, telegramId: {existing_user.user_id}")
# Update telegram ID if it was missing in panel data but we have local user
if (
telegram_id_from_panel
and existing_user.user_id != telegram_id_from_panel
):
logging.warning(
f"TelegramId mismatch: panel={telegram_id_from_panel}, local={existing_user.user_id}"
)
if telegram_id_from_panel and existing_user.user_id != telegram_id_from_panel:
logging.warning(f"TelegramId mismatch: panel={telegram_id_from_panel}, local={existing_user.user_id}")
if not existing_user:
users_not_found_in_db += 1
if telegram_id_from_panel:
@@ -128,36 +100,26 @@ async def perform_sync(
"user_id": telegram_id_from_panel,
"username": None, # Username will be updated when user interacts with bot
"first_name": None, # Panel doesn't provide this info
"last_name": None, # Panel doesn't provide this info
"last_name": None, # Panel doesn't provide this info
"language_code": "ru", # Default language
"panel_user_uuid": panel_uuid,
"is_banned": False,
"referred_by_id": None,
"referred_by_id": None
}
new_user, was_created = await user_dal.create_user(
session, user_data
)
new_user, was_created = await user_dal.create_user(session, user_data)
if was_created:
users_created += 1
logging.info(
f"Created new user {telegram_id_from_panel} from panel sync with UUID {panel_uuid}"
)
logging.info(f"Created new user {telegram_id_from_panel} from panel sync with UUID {panel_uuid}")
existing_user = new_user
except Exception as e_create:
sync_errors.append(
f"Error creating user {telegram_id_from_panel}: {str(e_create)}"
)
logging.error(
f"Error creating user {telegram_id_from_panel}: {e_create}"
)
sync_errors.append(f"Error creating user {telegram_id_from_panel}: {str(e_create)}")
logging.error(f"Error creating user {telegram_id_from_panel}: {e_create}")
continue
else:
logging.debug(
f"Panel user with UUID {panel_uuid} (no telegramId) not found in local DB - skipping"
)
logging.debug(f"Panel user with UUID {panel_uuid} (no telegramId) not found in local DB - skipping")
continue
# User found in local DB
@@ -172,29 +134,20 @@ async def perform_sync(
existing_user.panel_user_uuid = panel_uuid
user_was_updated = True
users_uuid_updated += 1
logging.info(
f"Updated panel UUID for user {actual_user_id}: {panel_uuid}"
)
logging.info(f"Updated panel UUID for user {actual_user_id}: {panel_uuid}")
# Ensure panel description contains Telegram fields
try:
if panel_uuid and existing_user:
description_text = "\n".join(
[
existing_user.username or "",
existing_user.first_name or "",
existing_user.last_name or "",
]
)
description_text = "\n".join([
existing_user.username or "",
existing_user.first_name or "",
existing_user.last_name or "",
])
# Update description only when it differs from the current one on panel
current_panel_description = (
panel_user_dict.get("description") or ""
).strip()
current_panel_description = (panel_user_dict.get("description") or "").strip()
desired_description = description_text.strip()
if (
desired_description
and desired_description != current_panel_description
):
if desired_description and desired_description != current_panel_description:
await panel_service.update_user_details_on_panel(
panel_uuid, {"description": description_text}
)
@@ -206,13 +159,13 @@ async def perform_sync(
# Sync subscription data
panel_expire_at_iso = panel_user_dict.get("expireAt")
panel_status = panel_user_dict.get("status", "UNKNOWN")
if panel_expire_at_iso:
try:
panel_expire_at = datetime.fromisoformat(
panel_expire_at_iso.replace("Z", "+00:00")
)
# Prefer syncing by concrete subscription UUID (shortUuid/subscriptionUuid)
subscription_uuid_from_panel = (
panel_user_dict.get("subscriptionUuid")
@@ -220,27 +173,6 @@ async def perform_sync(
)
if subscription_uuid_from_panel:
# Если панель говорит, что подписка ACTIVE — сначала деактивируем все другие активные
if panel_status == "ACTIVE":
await session.execute(
update(Subscription)
.where(
Subscription.panel_user_uuid == panel_uuid,
Subscription.is_active.is_(True),
or_(
Subscription.panel_subscription_uuid
!= subscription_uuid_from_panel,
Subscription.panel_subscription_uuid.is_(
None
),
),
)
.values(
is_active=False,
status_from_panel="INACTIVE",
)
)
# Try to find subscription by its panel_subscription_uuid first (idempotent)
existing_sub_by_uuid = (
await subscription_dal.get_subscription_by_panel_subscription_uuid(
@@ -265,8 +197,7 @@ async def perform_sync(
subscriptions_updated += 1
user_was_updated = True
logging.info(
f"Synced existing subscription {existing_sub_by_uuid.subscription_id} "
f"for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
f"Synced existing subscription {existing_sub_by_uuid.subscription_id} for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
)
else:
# Create a new subscription only when we have a concrete subscription UUID
@@ -290,15 +221,12 @@ async def perform_sync(
subscriptions_created += 1
user_was_updated = True
logging.info(
f"Created subscription {created_sub.subscription_id} "
f"for user {actual_user_id} by panel_sub_uuid {subscription_uuid_from_panel}"
f"Created subscription {created_sub.subscription_id} for user {actual_user_id} by panel_sub_uuid {subscription_uuid_from_panel}"
)
else:
# No subscription UUID from panel: only update an already active subscription for this user/panel UUID
active_sub = (
await subscription_dal.get_active_subscription_by_user_id(
session, actual_user_id, panel_uuid
)
active_sub = await subscription_dal.get_active_subscription_by_user_id(
session, actual_user_id, panel_uuid
)
if active_sub:
await subscription_dal.update_subscription(
@@ -314,30 +242,23 @@ async def perform_sync(
subscriptions_updated += 1
user_was_updated = True
logging.info(
f"Updated active subscription {active_sub.subscription_id} "
f"for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
f"Updated active subscription {active_sub.subscription_id} for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
)
else:
# Without a concrete subscription UUID we avoid creating new records to keep sync idempotent
logging.debug(
f"No subscriptionUuid for panel user {panel_uuid}; skipped creation for user {actual_user_id}"
)
except Exception as e:
sync_errors.append(
f"Error syncing subscription for user {actual_user_id}: {str(e)}"
)
logging.error(
f"Error syncing subscription for user {actual_user_id}: {e}"
)
sync_errors.append(f"Error syncing subscription for user {actual_user_id}: {str(e)}")
logging.error(f"Error syncing subscription for user {actual_user_id}: {e}")
if user_was_updated:
users_updated += 1
except Exception as e_user:
sync_errors.append(
f"Error processing panel user {panel_user_dict.get('uuid', 'unknown')}: {str(e_user)}"
)
sync_errors.append(f"Error processing panel user {panel_user_dict.get('uuid', 'unknown')}: {str(e_user)}")
logging.error(f"Error syncing user: {e_user}")
# Update sync status
@@ -346,26 +267,14 @@ async def perform_sync(
default_lang = settings.DEFAULT_LANGUAGE
additional_stats = ""
if users_without_telegram_id > 0:
additional_stats += i18n_instance.gettext(
default_lang,
"admin_sync_no_telegram_id",
count=users_without_telegram_id,
)
additional_stats += i18n_instance.gettext(default_lang, "admin_sync_no_telegram_id", count=users_without_telegram_id)
if users_not_found_in_db > 0:
additional_stats += i18n_instance.gettext(
default_lang,
"admin_sync_not_found_in_db",
count=users_not_found_in_db,
)
additional_stats += i18n_instance.gettext(default_lang, "admin_sync_not_found_in_db", count=users_not_found_in_db)
if sync_errors:
additional_stats += i18n_instance.gettext(
default_lang, "admin_sync_errors", count=len(sync_errors)
)
additional_stats += i18n_instance.gettext(default_lang, "admin_sync_errors", count=len(sync_errors))
# Build full details using localization
details = i18n_instance.gettext(
default_lang,
"admin_sync_details",
details = i18n_instance.gettext(default_lang, "admin_sync_details",
panel_records_checked=panel_records_checked,
users_found_in_db=users_found_in_db,
users_created=users_created,
@@ -373,15 +282,11 @@ async def perform_sync(
subscriptions_synced_count=subscriptions_synced_count,
subscriptions_created=subscriptions_created,
subscriptions_updated=subscriptions_updated,
additional_stats=additional_stats,
additional_stats=additional_stats
)
await panel_sync_dal.update_panel_sync_status(
session,
status,
details,
panel_records_checked,
subscriptions_synced_count,
session, status, details, panel_records_checked, subscriptions_synced_count
)
await session.commit()
@@ -406,27 +311,19 @@ async def perform_sync(
"users_synced": users_found_in_db,
"users_created": users_created,
"subs_synced": subscriptions_synced_count,
"errors": sync_errors,
"errors": sync_errors
}
except Exception as e_sync_global:
await session.rollback()
logging.error(f"Global error during sync: {e_sync_global}", exc_info=True)
error_detail = f"Unexpected error during sync: {str(e_sync_global)}"
await panel_sync_dal.update_panel_sync_status(
session,
"failed",
error_detail,
panel_records_checked,
subscriptions_synced_count,
session, "failed", error_detail, panel_records_checked, subscriptions_synced_count
)
return {
"status": "failed",
"details": error_detail,
"errors": [str(e_sync_global)],
}
return {"status": "failed", "details": error_detail, "errors": [str(e_sync_global)]}
@router.message(Command("sync"))
@@ -469,40 +366,34 @@ async def sync_command_handler(
# Use the extracted perform_sync function
try:
sync_result = await perform_sync(panel_service, session, settings, i18n)
status = sync_result.get("status")
details = sync_result.get("details", "No details available")
errors = sync_result.get("errors", [])
# Simple confirmation message to admin
if status == "failed":
await bot.send_message(target_chat_id, _("sync_failed_simple"))
elif status == "completed_with_errors":
await bot.send_message(
target_chat_id,
_("sync_errors_simple", errors_count=len(errors)),
)
await bot.send_message(target_chat_id, _("sync_errors_simple", errors_count=len(errors)))
else:
await bot.send_message(target_chat_id, _("sync_success_simple"))
# Send notification to log channel with proper thread handling
try:
notification_service = NotificationService(bot, settings, i18n)
await notification_service.notify_panel_sync(
status,
details,
status, details,
sync_result.get("users_processed", 0),
sync_result.get("subs_synced", 0),
sync_result.get("subs_synced", 0)
)
except Exception as e_notification:
logging.error(f"Failed to send sync notification: {e_notification}")
except Exception as e_sync_global:
logging.error(
f"Global error during /sync command: {e_sync_global}", exc_info=True
)
logging.error(f"Global error during /sync command: {e_sync_global}", exc_info=True)
await bot.send_message(target_chat_id, _("sync_critical_error"))
# Send notification to log channel about failure
try:
notification_service = NotificationService(bot, settings, i18n)
@@ -510,9 +401,7 @@ async def sync_command_handler(
"failed", str(e_sync_global), 0, 0
)
except Exception as e_notification:
logging.error(
f"Failed to send sync failure notification: {e_notification}"
)
logging.error(f"Failed to send sync failure notification: {e_notification}")
@router.message(Command("syncstatus"))
@@ -531,9 +420,7 @@ async def sync_status_command_handler(
if status_record_model:
last_time_val = status_record_model.last_sync_time
last_time_str = (
last_time_val.strftime("%Y-%m-%d %H:%M:%S UTC")
if last_time_val
else "N/A"
last_time_val.strftime("%Y-%m-%d %H:%M:%S UTC") if last_time_val else "N/A"
)
details_val = status_record_model.details