"""
Планировщик задач (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'📥 Еженедельный экспорт завершён\n\n'
f'📄 Файл: {result.get("filename")}\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'❌ Ошибка экспорта\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:
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 = "📢 Объявление от администрации\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'✅ Пост опубликован!\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
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}')