470 lines
24 KiB
Python
470 lines
24 KiB
Python
"""
|
||
Планировщик задач (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('✅ Добавлена задача проверки запланированных постов')
|
||
|
||
# Еженедельная генерация дайджеста (воскресенье 15:00 UTC = 19:00 Ульяновск)
|
||
self.scheduler.add_job(
|
||
self.weekly_digest_generation,
|
||
CronTrigger.from_crontab('0 15 * * 0'), # Каждое воскресенье в 15:00 UTC (19:00 local)
|
||
id='weekly_digest',
|
||
name='Генерация еженедельного дайджеста',
|
||
replace_existing=True
|
||
)
|
||
logger.info('✅ Добавлена задача генерации дайджеста (воскресенье 19:00 Ульяновск)')
|
||
|
||
def stop(self):
|
||
"""Остановка планировщика"""
|
||
self.scheduler.shutdown()
|
||
logger.info('🛑 Планировщик остановлен')
|
||
|
||
async def weekly_export(self):
|
||
"""Еженедельный экспорт данных (Новая логика v2026)"""
|
||
try:
|
||
from database.db import AsyncSessionLocal
|
||
from services.chat_exporter import ChatExporter
|
||
|
||
async with AsyncSessionLocal() as session:
|
||
exporter = ChatExporter(session)
|
||
result = await exporter.run_full_export()
|
||
|
||
status_git = "✅ В Gitea" if result.get('git') else "❌ Ошибка Gitea"
|
||
|
||
await self.bot.send_message(
|
||
ADMIN_USER_ID,
|
||
f'📥 <b>Еженедельный экспорт завершён</b>\n\n'
|
||
f'📄 Файл: <code>{result.get("filename")}</code>\n'
|
||
f'📊 Сообщений: {result.get("messages")}\n'
|
||
f'🌐 Статус: {status_git}\n'
|
||
f'🕒 Пояс: Ульяновск (UTC+4)'
|
||
)
|
||
|
||
logger.info(f"✅ Weekly export finished: {result.get('filename')}")
|
||
|
||
except Exception as e:
|
||
logger.error(f'Ошибка экспорта: {e}')
|
||
await self.bot.send_message(ADMIN_USER_ID, f'❌ <b>Ошибка экспорта</b>\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'🗑️ <b>Пользователь удалён из чата</b>\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'⏰ <b>Напоминание о событии!</b>\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:
|
||
bot_token = config.BOT_TOKEN
|
||
base_url = f"https://api.telegram.org/bot{bot_token}"
|
||
proxy_url = config.get_proxy_url()
|
||
|
||
connector = None
|
||
if proxy_url:
|
||
if proxy_url.startswith('socks'):
|
||
from aiohttp_socks import ProxyConnector
|
||
connector = ProxyConnector.from_url(proxy_url)
|
||
else:
|
||
connector = aiohttp.TCPConnector()
|
||
|
||
async with aiohttp.ClientSession(connector=connector) as http_session:
|
||
# Функция для отправки контента (в чат или юзеру)
|
||
async def send_content(target_chat_id, thread_id=None, is_private=False):
|
||
prefix = "📢 <b>Объявление от администрации</b>\n\n" if is_private else ""
|
||
caption = f"{prefix}{post.text}"
|
||
|
||
# СЛУЧАЙ 1: И ФОТО, И ДОКУМЕНТ -> MediaGroup
|
||
if post.photo_file_id and post.document_file_id:
|
||
media = [
|
||
{"type": "photo", "media": post.photo_file_id},
|
||
{"type": "document", "media": post.document_file_id, "caption": caption, "parse_mode": "HTML"}
|
||
]
|
||
url = f"{base_url}/sendMediaGroup"
|
||
payload = {"chat_id": target_chat_id, "media": json.dumps(media)}
|
||
if thread_id: payload["message_thread_id"] = thread_id
|
||
async with http_session.post(url, json=payload) as r:
|
||
return (await r.json())['result'][0]['message_id'] if r.status == 200 else None
|
||
|
||
# СЛУЧАЙ 2: ТОЛЬКО ДОКУМЕНТ
|
||
elif post.document_file_id:
|
||
url = f"{base_url}/sendDocument"
|
||
payload = {"chat_id": target_chat_id, "document": post.document_file_id, "caption": caption, "parse_mode": "HTML"}
|
||
if thread_id: payload["message_thread_id"] = thread_id
|
||
async with http_session.post(url, json=payload) as r:
|
||
return (await r.json())['result']['message_id'] if r.status == 200 else None
|
||
|
||
# СЛУЧАЙ 3: ТОЛЬКО ФОТО
|
||
elif post.photo_file_id:
|
||
url = f"{base_url}/sendPhoto"
|
||
payload = {"chat_id": target_chat_id, "photo": post.photo_file_id, "caption": caption, "parse_mode": "HTML"}
|
||
if thread_id: payload["message_thread_id"] = thread_id
|
||
async with http_session.post(url, json=payload) as r:
|
||
return (await r.json())['result']['message_id'] if r.status == 200 else None
|
||
|
||
# СЛУЧАЙ 4: ТОЛЬКО ТЕКСТ
|
||
else:
|
||
url = f"{base_url}/sendMessage"
|
||
payload = {"chat_id": target_chat_id, "text": caption, "parse_mode": "HTML"}
|
||
if thread_id: payload["message_thread_id"] = thread_id
|
||
async with http_session.post(url, json=payload) as r:
|
||
return (await r.json())['result']['message_id'] if r.status == 200 else None
|
||
|
||
# 1. Отправка в основной чат
|
||
msg_id = await send_content(ADMIN_CHAT_ID, thread_id=post.topic_id)
|
||
|
||
# 2. Если нужно, рассылка жильцам в личку
|
||
if post.recipients in ['all_verified', 'all_and_chat']:
|
||
stmt_users = select(User).where(User.verified == True)
|
||
users = (await session.execute(stmt_users)).scalars().all()
|
||
for user in users:
|
||
try:
|
||
await send_content(user.user_id, is_private=True)
|
||
except Exception as e:
|
||
logger.error(f"Failed send to {user.user_id}: {e}")
|
||
|
||
# Финализация
|
||
post.status = 'sent'
|
||
post.sent_at = datetime.utcnow()
|
||
post.message_id = msg_id
|
||
await session.commit()
|
||
logger.info(f"✅ Пост {post.id} отправлен")
|
||
|
||
# Обновляем статус поста
|
||
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'✅ <b>Пост опубликован!</b>\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'❌ <b>Ошибка публикации поста!</b>\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
|
||
|
||
async def weekly_digest_generation(self):
|
||
"""Еженедельная генерация дайджеста (воскресенье 09:00)"""
|
||
try:
|
||
from handlers.digest import auto_generate_digest
|
||
await auto_generate_digest(self.bot)
|
||
logger.info('✅ Еженедельный дайджест сгенерирован')
|
||
except Exception as e:
|
||
logger.error(f'Ошибка генерации дайджеста: {e}')
|