""" Планировщик задач (cron) """ import logging from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select from datetime import datetime, timedelta from database.models import User from services.analytics import ChatAnalytics from services.spy_detector import SpyDetector from config import ADMIN_USER_ID, EXPORT_SCHEDULE, SPY_DETECTION_ENABLED import config logger = logging.getLogger(__name__) class Scheduler: """Планировщик фоновых задач""" def __init__(self, bot): self.bot = bot self.scheduler = AsyncIOScheduler() def start(self): """Запуск планировщика""" # Еженедельный экспорт self.scheduler.add_job( self.weekly_export, CronTrigger.from_crontab(EXPORT_SCHEDULE), id='weekly_export', name='Еженедельный экспорт чата', replace_existing=True ) # Ежедневная проверка шпионов if SPY_DETECTION_ENABLED: self.scheduler.add_job( self.daily_spy_check, CronTrigger.from_crontab('0 2 * * *'), # Каждый день в 02:00 id='daily_spy_check', name='Проверка шпионов', replace_existing=True ) # Ежедневная проверка неверифицированных (каждый день в 03:00) self.scheduler.add_job( self.daily_unverified_check, CronTrigger.from_crontab('0 3 * * *'), id='daily_unverified_check', name='Проверка неверифицированных', replace_existing=True ) # Ежечасное обновление активности self.scheduler.add_job( self.hourly_activity_update, CronTrigger.from_crontab('0 * * * *'), # Каждый час id='hourly_activity', name='Обновление активности', replace_existing=True ) # Ежечасная проверка напоминаний об отключениях self.scheduler.add_job( self.hourly_schedule_reminders, CronTrigger.from_crontab('0 * * * *'), # Каждый час id='hourly_schedule_reminders', name='Напоминания об отключениях', replace_existing=True ) # Ежедневная очистка просроченных объявлений (каждый день в 04:00) self.scheduler.add_job( self.daily_ads_cleanup, CronTrigger.from_crontab('0 4 * * *'), id='daily_ads_cleanup', name='Очистка просроченных объявлений', replace_existing=True ) # Ежедневная рассылка напоминаний о платежах (5 и 9 числа в 09:00) self.scheduler.add_job( self.daily_payment_reminders, CronTrigger.from_crontab('0 9 5,9 * *'), id='daily_payment_reminders', name='Напоминания о платежах ЖКХ', replace_existing=True ) # Ежечасная проверка напоминаний о событиях (каждый час) self.scheduler.add_job( self.hourly_event_reminders, CronTrigger.from_crontab('0 * * * *'), id='hourly_event_reminders', name='Напоминания о событиях', replace_existing=True ) self.scheduler.start() logger.info('✅ Планировщик запущен') # Добавляем задачу проверки запланированных постов (каждую минуту) self.scheduler.add_job( self.process_scheduled_posts, CronTrigger.from_crontab('* * * * *'), # Каждую минуту id='scheduled_posts_checker', name='Проверка запланированных постов', replace_existing=True ) logger.info('✅ Добавлена задача проверки запланированных постов') def stop(self): """Остановка планировщика""" self.scheduler.shutdown() logger.info('🛑 Планировщик остановлен') async def weekly_export(self): """Еженедельный экспорт данных""" try: from database.db import AsyncSessionLocal async with AsyncSessionLocal() as session: analytics = ChatAnalytics(session) filepath = await analytics.export_to_json() await self.bot.send_message( ADMIN_USER_ID, f'📥 Еженедельный экспорт выполнен\n\n' f'Файл сохранён: {filepath}\n' f'Время: {datetime.utcnow().strftime("%d.%m.%Y %H:%M")}' ) logger.info(f'Еженедельный экспорт: {filepath}') except Exception as e: logger.error(f'Ошибка экспорта: {e}') await self.bot.send_message( ADMIN_USER_ID, f'❌ Ошибка экспорта\n\n{str(e)}' ) async def daily_spy_check(self): """Ежедневная проверка на шпионов""" try: from database.db import AsyncSessionLocal async with AsyncSessionLocal() as session: spy_detector = SpyDetector(session) # Получаем всех подозрительных suspicious = await spy_detector.get_suspicious_users(min_score=30) for user in suspicious: # Обновляем score changed = await spy_detector.update_spy_score(user) # Если стал шпионом (score > 60) — уведомляем админа if user.spy_score > 60 and changed: from utils.formatters import format_spy_alert flags = user.spy_flags import json flags_list = json.loads(flags) if flags else [] alert_text = format_spy_alert(user, flags_list) await self.bot.send_message( ADMIN_USER_ID, alert_text, parse_mode='HTML' ) logger.warning(f'🚨 Обнаружен шпион: user_id={user.user_id}') except Exception as e: logger.error(f'Ошибка проверки шпионов: {e}') async def daily_unverified_check(self): """ Ежедневная проверка пользователей с отозванной верификацией Удаляет тех кто не верифицировался более 7 дней """ try: from database.db import AsyncSessionLocal from sqlalchemy import update async with AsyncSessionLocal() as session: # Находим пользователей с отозванной верификацией более 7 дней week_ago = datetime.utcnow() - timedelta(days=7) stmt = ( select(User) .where(User.verified == False) .where(User.verification_date <= week_ago) .where(User.is_banned == False) ) result = await session.execute(stmt) users = list(result.scalars().all()) if users: logger.info(f'Найдено {len(users)} пользователей для удаления') for user in users: # Пытаемся удалить из чата try: await self.bot.chat.ban(user.user_id) user.is_banned = True await session.commit() logger.info(f'Удалён из чата: user_id={user.user_id}') # Уведомляем админа await self.bot.send_message( ADMIN_USER_ID, f'🗑️ Пользователь удалён из чата\n\n' f'Пользователь: {user.full_name}\n' f'Квартира: {user.apartment}\n' f'Причина: Не прошёл верификацию за 7 дней', parse_mode='HTML' ) except Exception as e: logger.error(f'Не удалось удалить пользователя {user.user_id}: {e}') except Exception as e: logger.error(f'Ошибка проверки неверифицированных: {e}') async def hourly_activity_update(self): """Ежечасное обновление счётчиков активности""" try: from database.db import AsyncSessionLocal from sqlalchemy import update async with AsyncSessionLocal() as session: # Обновляем last_seen для всех пользователей # (в реальности нужно отслеживать по сообщениям) pass except Exception as e: logger.error(f'Ошибка обновления активности: {e}') async def hourly_schedule_reminders(self): """Ежечасная проверка напоминаний об отключениях""" try: from handlers.schedule import check_and_send_reminders await check_and_send_reminders(self.bot) except Exception as e: logger.error(f'Ошибка напоминаний об отключениях: {e}') async def daily_ads_cleanup(self): """Ежедневная очистка просроченных объявлений""" try: from handlers.ads import cleanup_expired_ads await cleanup_expired_ads() except Exception as e: logger.error(f'Ошибка очистки объявлений: {e}') async def daily_payment_reminders(self): """Ежедневная рассылка напоминаний о платежах""" try: from handlers.payments import send_payment_reminders await send_payment_reminders(self.bot) except Exception as e: logger.error(f'Ошибка напоминаний о платежах: {e}') async def hourly_event_reminders(self): """Ежечасная проверка напоминаний о событиях""" try: from services.events import EventsService from database.db import AsyncSessionLocal from config import ADMIN_CHAT_ID async with AsyncSessionLocal() as session: events_service = EventsService(session) events = await events_service.get_events_for_reminder(hours_ahead=24) for event in events: # Отправляем напоминание в чат try: date_str = event['event_date'].strftime('%d.%m.%Y %H:%M') await self.bot.send_message( ADMIN_CHAT_ID, f'⏰ Напоминание о событии!\n\n' f'📅 {event["title"]}\n' f'🗓️ {date_str}\n' f'📍 {event["location"] or "Место не указано"}\n\n' f'Не забудьте принять участие!', parse_mode='HTML' ) await events_service.mark_reminder_sent(event['id']) logger.info(f'Напоминание о событии {event["id"]} отправлено') except Exception as e: logger.error(f'Не удалось отправить напоминание о событии: {e}') except Exception as e: logger.error(f'Ошибка напоминаний о событиях: {e}') async def process_scheduled_posts(self): """Проверка и отправка запланированных постов""" try: from database.db import AsyncSessionLocal from database.models import ScheduledPost, User from config import ADMIN_CHAT_ID, get_topic_id from sqlalchemy import update import aiohttp # Защита от дублирования — блокировка на уровне процесса if hasattr(self, '_scheduled_posts_lock') and self._scheduled_posts_lock: logger.warning("⚠️ Посты уже обрабатываются, пропускаю") return self._scheduled_posts_lock = True async with AsyncSessionLocal() as session: # Ищем посты которые пора отправить # SELECT FOR UPDATE — блокировка на уровне БД stmt = ( select(ScheduledPost) .where(ScheduledPost.status == 'pending') .where(ScheduledPost.scheduled_time <= datetime.utcnow()) ) result = await session.execute(stmt) posts_to_send = list(result.scalars().all()) if not posts_to_send: self._scheduled_posts_lock = False return logger.info(f"📅 Найдено {len(posts_to_send)} постов для отправки") # Сразу обновляем статус на 'sending' чтобы другие процессы не взяли for post in posts_to_send: post.status = 'sending' # Промежуточный статус await session.commit() for post in posts_to_send: try: # Отправляем пост message_id = None bot_token = config.BOT_TOKEN base_url = f"https://api.telegram.org/bot{bot_token}" # Настройка прокси для aiohttp proxy_url = config.get_proxy_url() connector = None if proxy_url: # Для socks прокси используем special connector if proxy_url.startswith('socks'): from aiohttp_socks import ProxyConnector connector = ProxyConnector.from_url(proxy_url) else: # HTTP прокси connector = aiohttp.TCPConnector() if post.photo_file_id: # С фото - отправляем по file_id url = f"{base_url}/sendPhoto" params = { 'chat_id': ADMIN_CHAT_ID, 'photo': post.photo_file_id, 'caption': post.text, 'parse_mode': 'HTML' } # Не добавляем message_thread_id для general (1) — это общий чат if post.topic_id and post.topic_id > 1: params['message_thread_id'] = post.topic_id async with aiohttp.ClientSession(connector=connector) as http_session: async with http_session.post(url, json=params, timeout=30) as resp: if resp.status == 200: result_data = await resp.json() message_id = result_data['result']['message_id'] else: error_text = await resp.text() raise Exception(f"Telegram API {resp.status}: {error_text}") else: # Только текст url = f"{base_url}/sendMessage" params = { 'chat_id': ADMIN_CHAT_ID, 'text': post.text, 'parse_mode': 'HTML' } # Не добавляем message_thread_id для general (1) — это общий чат if post.topic_id and post.topic_id > 1: params['message_thread_id'] = post.topic_id async with aiohttp.ClientSession(connector=connector) as http_session: async with http_session.post(url, json=params, timeout=30) as resp: if resp.status == 200: result_data = await resp.json() message_id = result_data['result']['message_id'] else: error_text = await resp.text() raise Exception(f"Telegram API {resp.status}: {error_text}") # Если нужно отправить в личку пользователям if post.recipients in ['all_verified', 'all_and_chat']: stmt_users = select(User).where(User.verified == True) result_users = await session.execute(stmt_users) users = list(result_users.scalars().all()) for user in users: try: if post.photo_file_id: url = f"{base_url}/sendPhoto" params = { 'chat_id': user.user_id, 'photo': post.photo_file_id, 'caption': f"📢 Объявление от администрации\n\n{post.text}", 'parse_mode': 'HTML' } else: url = f"{base_url}/sendMessage" params = { 'chat_id': user.user_id, 'text': f"📢 Объявление от администрации\n\n{post.text}", 'parse_mode': 'HTML' } async with aiohttp.ClientSession(connector=connector) as http_session: async with http_session.post(url, json=params, timeout=10) as resp: if resp.status != 200: logger.error(f"Не удалось отправить пользователю {user.user_id}") except Exception as e: logger.error(f"Ошибка отправки пользователю {user.user_id}: {e}") # Обновляем статус поста post.status = 'sent' post.sent_at = datetime.utcnow() post.message_id = message_id await session.commit() logger.info(f"✅ Пост отправлен: id={post.id}, тема={post.topic_name}") # Уведомляем админа await self.bot.send_message( config.ADMIN_USER_ID, f'✅ Пост опубликован!\n\n' f'ID: {post.id}\n' f'Тема: {post.get_topic_emoji()} {post.topic_name}\n' f'Время: {post.scheduled_time.strftime("%d.%m.%Y %H:%M")}\n' f'Статус: {post.get_status_emoji()} {post.status}', parse_mode='HTML' ) except Exception as e: # Ошибка при отправке post.status = 'failed' post.error_message = str(e)[:500] await session.commit() logger.error(f"❌ Ошибка отправки поста {post.id}: {e}") # Уведомляем админа об ошибке await self.bot.send_message( config.ADMIN_USER_ID, f'❌ Ошибка публикации поста!\n\n' f'ID: {post.id}\n' f'Тема: {post.topic_name}\n' f'Ошибка: {str(e)[:200]}', parse_mode='HTML' ) except Exception as e: logger.error(f'Ошибка обработки запланированных постов: {e}') finally: # Снимаем блокировку self._scheduled_posts_lock = False