Auto block-unblock chat member / Backup S3
This commit is contained in:
Binary file not shown.
@@ -72,6 +72,10 @@ AUDIT_REDIS_IDENTITY_PREFIX = "audit:user:identity:"
|
||||
AUDIT_REDIS_USER_TTL_SEC = 25 * 3600
|
||||
AUDIT_REDIS_DRAIN_BATCH = 1000
|
||||
|
||||
BLOCKED_EVENTS_REDIS_KEY = "blocked:events"
|
||||
BLOCKED_DRAIN_BATCH = 500
|
||||
BLOCKED_DRAIN_INTERVAL_SEC = 30
|
||||
|
||||
ERROR_THROTTLE_WINDOW_SEC = 60
|
||||
ERROR_THROTTLE_MAX_KEYS = 500
|
||||
ERROR_THROTTLE_MESSAGE_MAX_LEN = 120
|
||||
|
||||
@@ -55,6 +55,40 @@ def backup_thread_loop(stop_event, _bot, _sessionmaker) -> None:
|
||||
loop.close()
|
||||
|
||||
|
||||
async def blocked_drain_loop(_bot, sessionmaker) -> None:
|
||||
from core.cache_config import BLOCKED_DRAIN_BATCH, BLOCKED_DRAIN_INTERVAL_SEC, BLOCKED_EVENTS_REDIS_KEY
|
||||
from core.redis_cache import cache_lpop_batch
|
||||
from database.bans import remove_blocked_user_ids, save_blocked_user_ids
|
||||
|
||||
while True:
|
||||
try:
|
||||
events = await cache_lpop_batch(BLOCKED_EVENTS_REDIS_KEY, BLOCKED_DRAIN_BATCH)
|
||||
if events:
|
||||
final: dict[int, str] = {}
|
||||
for ev in events:
|
||||
tg_id = ev.get("tg_id")
|
||||
action = ev.get("action")
|
||||
if tg_id and action:
|
||||
final[int(tg_id)] = action
|
||||
|
||||
to_add = [tid for tid, act in final.items() if act == "block"]
|
||||
to_remove = [tid for tid, act in final.items() if act == "unblock"]
|
||||
|
||||
if to_add or to_remove:
|
||||
async with sessionmaker() as session:
|
||||
if to_add:
|
||||
await save_blocked_user_ids(session, to_add)
|
||||
if to_remove:
|
||||
await remove_blocked_user_ids(session, to_remove)
|
||||
await session.commit()
|
||||
logger.info(
|
||||
"[BlockedDrain] add={}, remove={}", len(to_add), len(to_remove)
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error("[BlockedDrain] Ошибка: {}", e)
|
||||
await asyncio.sleep(BLOCKED_DRAIN_INTERVAL_SEC)
|
||||
|
||||
|
||||
async def server_checks_loop(_bot, sessionmaker) -> None:
|
||||
from config import PING_TIME
|
||||
from servers import check_servers
|
||||
|
||||
@@ -15,6 +15,7 @@ from core.tasks.cron_tasks import (
|
||||
from core.tasks.loop_tasks import (
|
||||
backup_loop,
|
||||
backup_thread_loop,
|
||||
blocked_drain_loop,
|
||||
notifications_loop,
|
||||
scheduled_broadcasts_loop_task,
|
||||
server_checks_loop,
|
||||
@@ -51,6 +52,8 @@ def register_periodic_tasks() -> None:
|
||||
else:
|
||||
periodic_task_manager.register_thread_loop_task("backup", backup_thread_loop)
|
||||
|
||||
periodic_task_manager.register_loop_task("blocked_drain", blocked_drain_loop)
|
||||
|
||||
if process_budget > 0:
|
||||
periodic_task_manager.register_process_loop_task("server_checks", server_checks_loop)
|
||||
process_budget -= 1
|
||||
|
||||
+15
-1
@@ -1,4 +1,4 @@
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy import delete, select
|
||||
from sqlalchemy.dialects.postgresql import insert
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
@@ -37,3 +37,17 @@ async def save_blocked_user_ids(session: AsyncSession, tg_ids: list[int]) -> Non
|
||||
await session.execute(stmt)
|
||||
total += len(values)
|
||||
logger.info(f"📝 Добавлено до {total} пользователей в blocked_users")
|
||||
|
||||
|
||||
async def remove_blocked_user_ids(session: AsyncSession, tg_ids: list[int]) -> None:
|
||||
"""Удаление списка telegram id из таблицы blocked_users батчами по 500."""
|
||||
if not tg_ids:
|
||||
return
|
||||
batch_size = 500
|
||||
total = 0
|
||||
for i in range(0, len(tg_ids), batch_size):
|
||||
batch = tg_ids[i : i + batch_size]
|
||||
stmt = delete(BlockedUser).where(BlockedUser.tg_id.in_(batch))
|
||||
result = await session.execute(stmt)
|
||||
total += result.rowcount
|
||||
logger.info(f"🗑 Удалено {total} пользователей из blocked_users")
|
||||
|
||||
@@ -4,6 +4,7 @@ from aiogram import Router
|
||||
|
||||
from .admin import router as admin_router
|
||||
from .captcha import router as captcha_router
|
||||
from .chat_member import router as chat_member_router
|
||||
from .coupons import router as coupons_router
|
||||
from .donate import router as donate_router
|
||||
from .instructions import router as instructions_router
|
||||
@@ -19,6 +20,7 @@ from .tariffs import router as tariff_router
|
||||
router = Router(name="handlers_main_router")
|
||||
|
||||
router.include_routers(
|
||||
chat_member_router,
|
||||
start_router,
|
||||
captcha_router,
|
||||
profile_router,
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
from aiogram import Router
|
||||
from aiogram.types import ChatMemberUpdated
|
||||
|
||||
from core.cache_config import BLOCKED_EVENTS_REDIS_KEY
|
||||
from core.redis_cache import cache_rpush
|
||||
from logger import logger
|
||||
|
||||
router = Router(name="chat_member_router")
|
||||
|
||||
|
||||
@router.my_chat_member()
|
||||
async def on_my_chat_member(event: ChatMemberUpdated):
|
||||
tg_id = event.from_user.id
|
||||
new_status = event.new_chat_member.status
|
||||
|
||||
if new_status == "kicked":
|
||||
await cache_rpush(BLOCKED_EVENTS_REDIS_KEY, {"tg_id": tg_id, "action": "block"})
|
||||
logger.info(f"[ChatMember] Пользователь {tg_id} заблокировал бота → событие в Redis")
|
||||
|
||||
elif new_status == "member":
|
||||
old_status = event.old_chat_member.status if event.old_chat_member else None
|
||||
if old_status in ("kicked", "left"):
|
||||
await cache_rpush(BLOCKED_EVENTS_REDIS_KEY, {"tg_id": tg_id, "action": "unblock"})
|
||||
logger.info(f"[ChatMember] Пользователь {tg_id} разблокировал бота → событие в Redis")
|
||||
@@ -271,7 +271,7 @@ async def back_to_currency(callback_query: CallbackQuery, state: FSMContext, ses
|
||||
|
||||
@router.callback_query(F.data == "back_to_pay")
|
||||
async def back_to_pay(callback_query: CallbackQuery, state: FSMContext, session: AsyncSession):
|
||||
return await balance_handler(callback_query, session)
|
||||
return await balance_handler(callback_query, state, session)
|
||||
|
||||
|
||||
@router.callback_query(F.data == "pay_tribute")
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -13,6 +13,7 @@ async-timeout==4.0.3
|
||||
asyncpg==0.30.0
|
||||
attrs==24.2.0
|
||||
babel==2.17.0
|
||||
boto3>=1.35.0
|
||||
cachetools==5.5.1
|
||||
certifi==2023.11.17
|
||||
cffi==1.17.1
|
||||
|
||||
+115
-102
@@ -14,15 +14,20 @@ from aiogram.types import BufferedInputFile
|
||||
from config import (
|
||||
ADMIN_ID,
|
||||
BACKUP_CAPTION,
|
||||
BACKUP_CHANNEL_ID,
|
||||
BACKUP_CHANNEL_THREAD_ID,
|
||||
BACKUP_CREATE_ARCHIVE,
|
||||
BACKUP_DESTINATION,
|
||||
BACKUP_INCLUDE_CONFIG,
|
||||
BACKUP_INCLUDE_DB,
|
||||
BACKUP_INCLUDE_IMG,
|
||||
BACKUP_INCLUDE_TEXTS,
|
||||
BACKUP_OTHER_BOT_TOKEN,
|
||||
BACKUP_SEND_MODE,
|
||||
BACKUP_S3_ACCESS_KEY,
|
||||
BACKUP_S3_BUCKET,
|
||||
BACKUP_S3_ENDPOINT,
|
||||
BACKUP_S3_KEEP,
|
||||
BACKUP_S3_PATH,
|
||||
BACKUP_S3_REGION,
|
||||
BACKUP_S3_SECRET_KEY,
|
||||
BACK_DIR,
|
||||
DB_NAME,
|
||||
DB_PASSWORD,
|
||||
@@ -37,6 +42,21 @@ from logger import logger
|
||||
DOCKER_POSTGRES_CONTAINER = "solobot-postgres"
|
||||
|
||||
|
||||
def _s3_configured() -> bool:
|
||||
return bool(BACKUP_S3_ENDPOINT and BACKUP_S3_ACCESS_KEY and BACKUP_S3_SECRET_KEY and BACKUP_S3_BUCKET)
|
||||
|
||||
|
||||
def _parse_destination() -> tuple[str | None, int | None]:
|
||||
"""Парсит BACKUP_DESTINATION в (chat_id, thread_id)."""
|
||||
raw = BACKUP_DESTINATION.strip()
|
||||
if not raw:
|
||||
return None, None
|
||||
parts = raw.split(":", 1)
|
||||
chat_id = parts[0]
|
||||
thread_id = int(parts[1]) if len(parts) > 1 and parts[1].strip() else None
|
||||
return chat_id, thread_id
|
||||
|
||||
|
||||
def _find_docker_postgres_container() -> str | None:
|
||||
if shutil.which("docker") is None:
|
||||
return None
|
||||
@@ -90,11 +110,8 @@ def _create_database_backup_via_docker(filename: Path, container: str) -> None:
|
||||
|
||||
async def backup_database(bot_instance: Bot | None = None) -> Exception | None:
|
||||
"""
|
||||
Создает резервную копию базы данных (или полный архив) и отправляет его администраторам.
|
||||
Блокирующие операции (pg_dump и т.д.) выполняются в пуле процессов, не блокируя event loop и используя другие ядра CPU.
|
||||
|
||||
Returns:
|
||||
Optional[Exception]: Исключение в случае ошибки или None при успешном выполнении
|
||||
Создает резервную копию и отправляет в S3 или Telegram.
|
||||
Блокирующие операции выполняются в пуле потоков/процессов.
|
||||
"""
|
||||
from core.executor import run_io
|
||||
|
||||
@@ -112,9 +129,15 @@ async def backup_database(bot_instance: Bot | None = None) -> Exception | None:
|
||||
|
||||
logger.info("[Backup] Файл создан: {}", backup_file_path)
|
||||
try:
|
||||
await _send_backup_to_admins(backup_file_path, bot_instance=bot_instance)
|
||||
exception = await run_io(_cleanup_old_backups)
|
||||
if _s3_configured():
|
||||
s3_err = await run_io(_upload_to_s3, backup_file_path)
|
||||
if s3_err:
|
||||
logger.error("[Backup] Ошибка S3: {}", s3_err)
|
||||
return s3_err
|
||||
else:
|
||||
await _send_backup_telegram(backup_file_path, bot_instance=bot_instance)
|
||||
|
||||
exception = await run_io(_cleanup_old_backups)
|
||||
if exception:
|
||||
logger.error("[Backup] Ошибка при очистке старых: {}", exception)
|
||||
return exception
|
||||
@@ -126,12 +149,6 @@ async def backup_database(bot_instance: Bot | None = None) -> Exception | None:
|
||||
|
||||
|
||||
def _create_database_backup() -> tuple[str | None, Exception | None]:
|
||||
"""
|
||||
Создает резервную копию базы данных PostgreSQL.
|
||||
|
||||
Returns:
|
||||
Tuple[Optional[str], Optional[Exception]]: Путь к файлу бэкапа и исключение (если произошла ошибка)
|
||||
"""
|
||||
date_formatted = datetime.now().strftime("%Y-%m-%d-%H%M%S")
|
||||
pid_suffix = os.getpid()
|
||||
|
||||
@@ -183,12 +200,6 @@ def _create_database_backup() -> tuple[str | None, Exception | None]:
|
||||
|
||||
|
||||
def _create_backup_archive() -> tuple[str | None, Exception | None]:
|
||||
"""
|
||||
Создает архив (.tar.gz) с выбранными компонентами бекапа.
|
||||
|
||||
Returns:
|
||||
Tuple[Optional[str], Optional[Exception]]: Путь к файлу архива и исключение (если произошла ошибка)
|
||||
"""
|
||||
date_formatted = datetime.now().strftime("%Y-%m-%d-%H%M%S")
|
||||
pid_suffix = os.getpid()
|
||||
backup_dir = Path(BACK_DIR)
|
||||
@@ -252,12 +263,6 @@ def _create_backup_archive() -> tuple[str | None, Exception | None]:
|
||||
|
||||
|
||||
def _cleanup_old_backups() -> Exception | None:
|
||||
"""
|
||||
Удаляет бэкапы старше 3 дней (как .sql, так и .tar.gz файлы).
|
||||
|
||||
Returns:
|
||||
Optional[Exception]: Исключение в случае ошибки или None при успешном выполнении
|
||||
"""
|
||||
try:
|
||||
backup_dir = Path(BACK_DIR)
|
||||
if not backup_dir.exists():
|
||||
@@ -286,94 +291,102 @@ def _cleanup_old_backups() -> Exception | None:
|
||||
return e
|
||||
|
||||
|
||||
async def create_backup_and_send_to_admins(client) -> None:
|
||||
"""
|
||||
Создает бэкап и отправляет администраторам через переданный клиент.
|
||||
def _create_s3_client():
|
||||
import boto3
|
||||
|
||||
Args:
|
||||
client: Клиент для работы с базой данных
|
||||
"""
|
||||
return boto3.client(
|
||||
"s3",
|
||||
endpoint_url=BACKUP_S3_ENDPOINT,
|
||||
aws_access_key_id=BACKUP_S3_ACCESS_KEY,
|
||||
aws_secret_access_key=BACKUP_S3_SECRET_KEY,
|
||||
region_name=BACKUP_S3_REGION or "us-east-1",
|
||||
)
|
||||
|
||||
|
||||
def _upload_to_s3(backup_file_path: str) -> Exception | None:
|
||||
"""Загружает бекап в S3 и чистит старые (синхронно, вызывается через run_io)."""
|
||||
try:
|
||||
s3 = _create_s3_client()
|
||||
prefix = BACKUP_S3_PATH.strip("/")
|
||||
object_key = f"{prefix}/{os.path.basename(backup_file_path)}"
|
||||
|
||||
s3.upload_file(backup_file_path, BACKUP_S3_BUCKET, object_key)
|
||||
logger.info("[Backup S3] Загружен: {}", object_key)
|
||||
|
||||
_cleanup_s3_backups(s3, prefix)
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.error("[Backup S3] Ошибка: {}", e)
|
||||
return e
|
||||
|
||||
|
||||
def _cleanup_s3_backups(s3, prefix: str) -> None:
|
||||
"""Удаляет старые бекапы в S3, оставляя BACKUP_S3_KEEP последних."""
|
||||
if BACKUP_S3_KEEP <= 0:
|
||||
return
|
||||
|
||||
all_objects = []
|
||||
paginator = s3.get_paginator("list_objects_v2")
|
||||
for page in paginator.paginate(Bucket=BACKUP_S3_BUCKET, Prefix=f"{prefix}/"):
|
||||
all_objects.extend(page.get("Contents", []))
|
||||
|
||||
all_objects.sort(key=lambda x: x["LastModified"])
|
||||
|
||||
if len(all_objects) <= BACKUP_S3_KEEP:
|
||||
return
|
||||
|
||||
to_delete = all_objects[: -BACKUP_S3_KEEP]
|
||||
for i in range(0, len(to_delete), 1000):
|
||||
batch = [{"Key": obj["Key"]} for obj in to_delete[i : i + 1000]]
|
||||
s3.delete_objects(Bucket=BACKUP_S3_BUCKET, Delete={"Objects": batch, "Quiet": True})
|
||||
|
||||
logger.info("[Backup S3] Удалено старых бекапов: {}", len(to_delete))
|
||||
|
||||
|
||||
async def create_backup_and_send_to_admins(client) -> None:
|
||||
await client.login()
|
||||
await client.database.export()
|
||||
|
||||
|
||||
async def _send_backup_to_admins(backup_file_path: str, bot_instance: Bot | None = None) -> None:
|
||||
"""
|
||||
Отправляет файл бэкапа всем администраторам через Telegram.
|
||||
|
||||
Args:
|
||||
backup_file_path: Путь к файлу бэкапа
|
||||
|
||||
Raises:
|
||||
Exception: При ошибке отправки файла
|
||||
"""
|
||||
async def _send_backup_telegram(backup_file_path: str, bot_instance: Bot | None = None) -> None:
|
||||
if not backup_file_path or not os.path.exists(backup_file_path):
|
||||
raise FileNotFoundError(f"Файл бэкапа не найден: {backup_file_path}")
|
||||
|
||||
active_bot = bot_instance
|
||||
if active_bot is None:
|
||||
own_session = False
|
||||
|
||||
if BACKUP_OTHER_BOT_TOKEN:
|
||||
active_bot = Bot(token=BACKUP_OTHER_BOT_TOKEN)
|
||||
own_session = True
|
||||
elif active_bot is None:
|
||||
from bot import bot as active_bot
|
||||
|
||||
async def send_default():
|
||||
for admin_id in ADMIN_ID:
|
||||
try:
|
||||
await active_bot.send_document(chat_id=admin_id, document=backup_input_file)
|
||||
logger.info("[Backup] Отправлено админу: {}", admin_id)
|
||||
except Exception as e:
|
||||
logger.error("[Backup] Не отправлено админу {}: {}", admin_id, e)
|
||||
|
||||
try:
|
||||
async with aiofiles.open(backup_file_path, "rb") as backup_file:
|
||||
backup_data = await backup_file.read()
|
||||
filename = os.path.basename(backup_file_path)
|
||||
backup_input_file = BufferedInputFile(file=backup_data, filename=filename)
|
||||
async with aiofiles.open(backup_file_path, "rb") as f:
|
||||
backup_data = await f.read()
|
||||
filename = os.path.basename(backup_file_path)
|
||||
backup_input_file = BufferedInputFile(file=backup_data, filename=filename)
|
||||
|
||||
if BACKUP_SEND_MODE == "default":
|
||||
await send_default()
|
||||
logger.info("[Backup] Отправлено всем админам")
|
||||
chat_id, thread_id = _parse_destination()
|
||||
|
||||
elif BACKUP_SEND_MODE == "channel":
|
||||
channel_id = BACKUP_CHANNEL_ID.strip()
|
||||
thread_id = BACKUP_CHANNEL_THREAD_ID.strip()
|
||||
if not channel_id:
|
||||
logger.error("[Backup] BACKUP_CHANNEL_ID не задан, fallback на default")
|
||||
await send_default()
|
||||
return
|
||||
send_kwargs = {"chat_id": channel_id, "document": backup_input_file}
|
||||
if thread_id:
|
||||
send_kwargs["message_thread_id"] = int(thread_id)
|
||||
if BACKUP_CAPTION:
|
||||
send_kwargs["caption"] = BACKUP_CAPTION
|
||||
if chat_id:
|
||||
send_kwargs: dict = {"chat_id": chat_id, "document": backup_input_file}
|
||||
if thread_id:
|
||||
send_kwargs["message_thread_id"] = thread_id
|
||||
if BACKUP_CAPTION:
|
||||
send_kwargs["caption"] = BACKUP_CAPTION
|
||||
await active_bot.send_document(**send_kwargs)
|
||||
logger.info("[Backup] Отправлено в {}{}", chat_id, f" (тред {thread_id})" if thread_id else "")
|
||||
else:
|
||||
for admin_id in ADMIN_ID:
|
||||
try:
|
||||
send_kwargs = {"chat_id": admin_id, "document": backup_input_file}
|
||||
if BACKUP_CAPTION:
|
||||
send_kwargs["caption"] = BACKUP_CAPTION
|
||||
await active_bot.send_document(**send_kwargs)
|
||||
logger.info("[Backup] Отправлено в канал: {} (топик: {})", channel_id, thread_id)
|
||||
logger.info("[Backup] Отправлено админу: {}", admin_id)
|
||||
except Exception as e:
|
||||
logger.error("[Backup] Не отправлено в канал {}: {}, fallback", channel_id, e)
|
||||
await send_default()
|
||||
|
||||
elif BACKUP_SEND_MODE == "bot":
|
||||
if not BACKUP_OTHER_BOT_TOKEN:
|
||||
logger.error("[Backup] BACKUP_OTHER_BOT_TOKEN не задан, fallback")
|
||||
await send_default()
|
||||
return
|
||||
other_bot = Bot(token=BACKUP_OTHER_BOT_TOKEN)
|
||||
try:
|
||||
for admin_id in ADMIN_ID:
|
||||
try:
|
||||
send_kwargs = {"chat_id": admin_id, "document": backup_input_file}
|
||||
if BACKUP_CAPTION:
|
||||
send_kwargs["caption"] = BACKUP_CAPTION
|
||||
await other_bot.send_document(**send_kwargs)
|
||||
logger.info("[Backup] Отправлено через другого бота админу: {}", admin_id)
|
||||
except Exception as e:
|
||||
logger.error("[Backup] Не отправлено админу {} через другого бота: {}", admin_id, e)
|
||||
await other_bot.session.close()
|
||||
except Exception as e:
|
||||
logger.error("[Backup] Ошибка через другого бота: {}, fallback", e)
|
||||
await send_default()
|
||||
else:
|
||||
logger.error("[Backup] Неизвестный BACKUP_SEND_MODE: {}, fallback", BACKUP_SEND_MODE)
|
||||
await send_default()
|
||||
except Exception as e:
|
||||
logger.error("[Backup] Ошибка при отправке: {}", e)
|
||||
raise
|
||||
logger.error("[Backup] Не отправлено админу {}: {}", admin_id, e)
|
||||
finally:
|
||||
if own_session:
|
||||
await active_bot.session.close()
|
||||
|
||||
Reference in New Issue
Block a user