All checks were successful
CI/CD Pipeline / build-and-deploy (push) Successful in 47s
Раньше о поломке сервиса админ узнавал только из алерта health-check или из ошибки живого пользователя. Теперь /stat отдаёт состояние всех пяти загрузчиков по запросу. У YouTube и Instagram дёргается /cookies/check, поэтому статусы различают причину: работает / cookies протухли / не извлекает видео (не cookies) / недоступен. Остальные три проверяются через /health. Проверка ходит в реальные YouTube и Instagram и занимает секунды, поэтому статистика отправляется сразу, а блок состояния дозаполняется правкой сообщения. Все сервисы опрашиваются параллельно — последовательно набежало бы под полминуты. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
410 lines
18 KiB
Python
410 lines
18 KiB
Python
"""
|
||
Admin Telegram Bot
|
||
Админский бот для получения статистики и всех скачанных видео
|
||
"""
|
||
import os
|
||
import asyncio
|
||
import logging
|
||
import sqlite3
|
||
from pathlib import Path
|
||
|
||
import httpx
|
||
from telegram import Bot, Update
|
||
from telegram.ext import Application, CommandHandler, ContextTypes, MessageHandler, filters
|
||
from telegram.request import HTTPXRequest
|
||
|
||
# Настройка логирования
|
||
logging.basicConfig(
|
||
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
||
level=logging.INFO
|
||
)
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Токен админ бота из переменных окружения
|
||
ADMIN_BOT_TOKEN = os.getenv('ADMIN_BOT_TOKEN')
|
||
|
||
# Токен клиентского бота — рассылка должна идти от его имени, иначе сообщения
|
||
# уходят в чат с админ-ботом, а не в чат с клиентским ботом, который видит пользователь
|
||
TELEGRAM_BOT_TOKEN = os.getenv('TELEGRAM_BOT_TOKEN')
|
||
|
||
# Базовая директория проекта
|
||
BASE_DIR = Path(__file__).resolve().parent
|
||
DATA_DIR = BASE_DIR / 'data'
|
||
DB_FILE = DATA_DIR / 'bot.db'
|
||
ADMIN_CHAT_ID_FILE = DATA_DIR / 'admin_chat_id.txt'
|
||
|
||
# Адреса загрузчиков — те же дефолты, что и в bot.py (все сервисы в network_mode: host).
|
||
# Флаг has_cookies: у YouTube и Instagram есть /cookies/check, у остальных только /health.
|
||
DOWNLOADER_SERVICES = [
|
||
('YouTube', os.getenv('YOUTUBE_DOWNLOADER_URL', 'http://localhost:5557'), True),
|
||
('Instagram', os.getenv('INSTAGRAM_DOWNLOADER_URL', 'http://localhost:5556'), True),
|
||
('VK', os.getenv('VK_DOWNLOADER_URL', 'http://localhost:5555'), False),
|
||
('Yapfiles', os.getenv('YAPFILES_DOWNLOADER_URL', 'http://localhost:5558'), False),
|
||
('TikTok', os.getenv('TIKTOK_DOWNLOADER_URL', 'http://localhost:5559'), False),
|
||
]
|
||
|
||
# /cookies/check ходит в реальный YouTube/Instagram, это несколько секунд
|
||
SERVICE_CHECK_TIMEOUT = httpx.Timeout(connect=5, read=45, write=10, pool=10)
|
||
|
||
# Расшифровка status из /cookies/check
|
||
COOKIE_STATUS_LABELS = {
|
||
'ok': '✅ работает',
|
||
'cookies_invalid': '🍪 cookies протухли',
|
||
'extraction_failed': '⚠️ не извлекает видео (не cookies)',
|
||
'no_cookies': '➖ без cookies',
|
||
}
|
||
|
||
|
||
async def check_service_status(name: str, base_url: str, has_cookies: bool) -> str:
|
||
"""Возвращает строку статуса одного загрузчика для /stat."""
|
||
try:
|
||
async with httpx.AsyncClient(timeout=SERVICE_CHECK_TIMEOUT) as client:
|
||
if has_cookies:
|
||
response = await client.post(f'{base_url}/cookies/check')
|
||
if response.status_code != 200:
|
||
return f"• {name}: ❌ HTTP {response.status_code}"
|
||
status = response.json().get('status', 'unknown')
|
||
label = COOKIE_STATUS_LABELS.get(status, f'⚠️ {status}')
|
||
return f"• {name}: {label}"
|
||
|
||
response = await client.get(f'{base_url}/health')
|
||
if response.status_code != 200:
|
||
return f"• {name}: ❌ HTTP {response.status_code}"
|
||
return f"• {name}: ✅ работает"
|
||
except Exception as e:
|
||
logger.warning(f"Проверка сервиса {name} не удалась: {e}")
|
||
return f"• {name}: ❌ недоступен"
|
||
|
||
|
||
async def get_services_status_text() -> str:
|
||
"""Опрашивает все загрузчики параллельно — последовательно это заняло бы десятки секунд."""
|
||
results = await asyncio.gather(*(
|
||
check_service_status(name, url, has_cookies)
|
||
for name, url, has_cookies in DOWNLOADER_SERVICES
|
||
))
|
||
return "\n".join(results)
|
||
|
||
|
||
def get_total_downloads() -> int:
|
||
"""Возвращает общее количество скачанных видео"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute('SELECT total_downloads FROM stats WHERE id = 1')
|
||
result = cursor.fetchone()
|
||
conn.close()
|
||
return result[0] if result else 0
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении количества скачанных видео: {e}")
|
||
return 0
|
||
|
||
|
||
def get_total_users() -> int:
|
||
"""Возвращает общее количество уникальных пользователей"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute('SELECT COUNT(*) FROM users')
|
||
result = cursor.fetchone()
|
||
conn.close()
|
||
return result[0] if result else 0
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении количества пользователей: {e}")
|
||
return 0
|
||
|
||
|
||
def get_users_sent_link_count() -> int:
|
||
"""Возвращает количество пользователей, отправивших хотя бы одну ссылку.
|
||
|
||
Ключ — реальный Telegram user_id (таблица user_links), а не chat_id: в
|
||
групповых чатах chat_id общий для всех участников и не годится для учёта
|
||
по отдельным людям.
|
||
"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute('SELECT COUNT(*) FROM user_links WHERE link_count > 0')
|
||
result = cursor.fetchone()
|
||
conn.close()
|
||
return result[0] if result else 0
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении количества активных пользователей: {e}")
|
||
return 0
|
||
|
||
|
||
def get_top_active_users(limit: int = 5) -> list[tuple[int, str, str, int]]:
|
||
"""Возвращает топ пользователей по количеству отправленных ссылок (по user_id)"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute(
|
||
'SELECT user_id, username, first_name, link_count FROM user_links '
|
||
'WHERE link_count > 0 ORDER BY link_count DESC LIMIT ?',
|
||
(limit,)
|
||
)
|
||
results = cursor.fetchall()
|
||
conn.close()
|
||
return results
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении топа активных пользователей: {e}")
|
||
return []
|
||
|
||
|
||
def get_error_stats() -> dict[str, int]:
|
||
"""Возвращает статистику ошибок по сервисам"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute('SELECT service, error_count FROM error_stats ORDER BY service')
|
||
results = cursor.fetchall()
|
||
conn.close()
|
||
return {service: count for service, count in results}
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении статистики ошибок: {e}")
|
||
return {}
|
||
|
||
|
||
def get_all_users() -> list[int]:
|
||
"""Возвращает chat_id всех пользователей бота"""
|
||
try:
|
||
conn = sqlite3.connect(str(DB_FILE))
|
||
cursor = conn.cursor()
|
||
cursor.execute('SELECT chat_id FROM users')
|
||
users = [row[0] for row in cursor.fetchall()]
|
||
conn.close()
|
||
return users
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при получении списка пользователей: {e}")
|
||
return []
|
||
|
||
|
||
def save_admin_chat_id(chat_id: int):
|
||
"""Сохраняет chat_id админа в файл"""
|
||
try:
|
||
DATA_DIR.mkdir(parents=True, exist_ok=True)
|
||
with open(ADMIN_CHAT_ID_FILE, 'w') as f:
|
||
f.write(str(chat_id))
|
||
logger.info(f"Сохранен chat_id админа: {chat_id}")
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при сохранении chat_id админа: {e}")
|
||
|
||
|
||
def get_admin_chat_id() -> int | None:
|
||
"""Получает сохраненный chat_id админа"""
|
||
try:
|
||
if ADMIN_CHAT_ID_FILE.exists():
|
||
with open(ADMIN_CHAT_ID_FILE, 'r') as f:
|
||
chat_id = f.read().strip()
|
||
if chat_id:
|
||
return int(chat_id)
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при чтении chat_id админа: {e}")
|
||
return None
|
||
|
||
|
||
async def do_broadcast(text: str) -> tuple[int, int, int]:
|
||
"""Рассылает текстовое сообщение всем пользователям клиентского бота.
|
||
|
||
Отправлять нужно от имени клиентского бота (TELEGRAM_BOT_TOKEN), а не
|
||
админского — пользователи переписываются именно с клиентским ботом, и
|
||
письмо от админ-бота уйдёт в отдельный, незнакомый им чат.
|
||
"""
|
||
users = get_all_users()
|
||
success = 0
|
||
failed = 0
|
||
|
||
async with Bot(token=TELEGRAM_BOT_TOKEN) as client_bot:
|
||
for chat_id in users:
|
||
try:
|
||
await client_bot.send_message(chat_id=chat_id, text=text)
|
||
success += 1
|
||
except Exception as e:
|
||
error_str = str(e).lower()
|
||
if 'blocked' not in error_str and 'chat not found' not in error_str:
|
||
logger.warning(f"Не удалось отправить сообщение {chat_id}: {e}")
|
||
failed += 1
|
||
|
||
await asyncio.sleep(0.05)
|
||
|
||
return success, failed, len(users)
|
||
|
||
|
||
async def start_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
|
||
"""Обрабатывает команду /start"""
|
||
chat_id = update.message.chat_id
|
||
saved_chat_id = get_admin_chat_id()
|
||
if saved_chat_id != chat_id:
|
||
save_admin_chat_id(chat_id)
|
||
if saved_chat_id is None:
|
||
await update.message.reply_text(
|
||
"✅ Админ бот активирован! Теперь вы будете получать все скачанные видео.\n\n"
|
||
"Доступные команды:\n"
|
||
"/stat — статистика бота\n"
|
||
"/send_messages — рассылка сообщения всем пользователям"
|
||
)
|
||
else:
|
||
await update.message.reply_text(
|
||
"Это админский бот.\n\n"
|
||
"Доступные команды:\n"
|
||
"/stat — статистика бота\n"
|
||
"/send_messages — рассылка сообщения всем пользователям"
|
||
)
|
||
else:
|
||
await update.message.reply_text(
|
||
"Доступные команды:\n"
|
||
"/stat — статистика бота\n"
|
||
"/send_messages — рассылка сообщения всем пользователям"
|
||
)
|
||
|
||
|
||
async def stat_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
|
||
"""Обрабатывает команду /stat"""
|
||
# Сохраняем chat_id админа при первом использовании
|
||
chat_id = update.message.chat_id
|
||
saved_chat_id = get_admin_chat_id()
|
||
if saved_chat_id != chat_id:
|
||
save_admin_chat_id(chat_id)
|
||
if saved_chat_id is None:
|
||
await update.message.reply_text("✅ Админ бот активирован! Теперь вы будете получать все скачанные видео.")
|
||
|
||
total_downloads = get_total_downloads()
|
||
total_users = get_total_users()
|
||
users_sent_link = get_users_sent_link_count()
|
||
top_users = get_top_active_users(10)
|
||
error_stats = get_error_stats()
|
||
|
||
# Форматируем статистику ошибок
|
||
service_names = {
|
||
'youtube': 'YouTube',
|
||
'instagram': 'Instagram',
|
||
'tiktok': 'TikTok',
|
||
'vk': 'VK',
|
||
'yapfiles': 'Yapfiles',
|
||
'unknown': 'Unknown'
|
||
}
|
||
|
||
error_lines = [
|
||
f"• {service_names.get(service, service)}: {count}"
|
||
for service, count in sorted(error_stats.items())
|
||
if count > 0
|
||
]
|
||
error_stats_text = "\n".join(error_lines) if error_lines else "Нет ошибок"
|
||
|
||
# Форматируем топ активных пользователей.
|
||
# (LRM) перед каждой строкой не даёт Telegram переразвернуть строку
|
||
# справа налево из-за RTL-символов в имени (арабский и т.п.) — иначе номер
|
||
# и счётчик "уезжают" вправо и ломают выравнивание списка.
|
||
top_lines = [
|
||
f"{i}. {('@' + username) if username else (first_name or str(uid))} — {link_count}"
|
||
for i, (uid, username, first_name, link_count) in enumerate(top_users, start=1)
|
||
]
|
||
top_users_text = "\n".join(top_lines) if top_lines else "Пока никто не отправлял ссылки"
|
||
|
||
stat_message = (
|
||
f"📊 Статистика бота:\n\n"
|
||
f"👥 Всего пользователей: {total_users}\n"
|
||
f"🔗 Отправляли ссылки: {users_sent_link}\n"
|
||
f"📹 Всего скачано видео: {total_downloads}\n\n"
|
||
f"🏆 Топ-10 активных пользователей:\n{top_users_text}\n\n"
|
||
f"❌ Ошибки по сервисам:\n{error_stats_text}"
|
||
)
|
||
|
||
# Проверка сервисов ходит в реальные YouTube/Instagram и занимает секунды,
|
||
# поэтому сначала отдаём статистику, а статус дозаполняем правкой сообщения
|
||
sent = await update.message.reply_text(
|
||
f"{stat_message}\n\n🔌 Состояние сервисов:\nпроверяю…"
|
||
)
|
||
|
||
try:
|
||
services_text = await get_services_status_text()
|
||
except Exception as e:
|
||
logger.error(f"Не удалось получить состояние сервисов: {e}")
|
||
services_text = "❌ не удалось проверить"
|
||
|
||
try:
|
||
await sent.edit_text(f"{stat_message}\n\n🔌 Состояние сервисов:\n{services_text}")
|
||
except Exception as e:
|
||
logger.warning(f"Не удалось обновить сообщение статистики: {e}")
|
||
|
||
|
||
async def send_messages_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
|
||
"""Обрабатывает команду /send_messages — запрашивает текст для рассылки"""
|
||
if not TELEGRAM_BOT_TOKEN:
|
||
await update.message.reply_text("❌ TELEGRAM_BOT_TOKEN не задан, рассылка недоступна")
|
||
return
|
||
context.user_data['awaiting_broadcast'] = True
|
||
await update.message.reply_text(
|
||
"Введите текст сообщения для рассылки всем пользователям.\n"
|
||
"Чтобы отменить — /cancel"
|
||
)
|
||
|
||
|
||
async def cancel_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
|
||
"""Обрабатывает команду /cancel — отменяет ожидание текста рассылки"""
|
||
if context.user_data.pop('awaiting_broadcast', None):
|
||
await update.message.reply_text("❌ Рассылка отменена")
|
||
else:
|
||
await update.message.reply_text("Нечего отменять")
|
||
|
||
|
||
async def handle_message(update: Update, context: ContextTypes.DEFAULT_TYPE):
|
||
"""Обрабатывает все сообщения (на случай если админ отправит что-то)"""
|
||
if context.user_data.pop('awaiting_broadcast', False):
|
||
text = f"📢 Service message\n\n{update.message.text}"
|
||
success, failed, total = await do_broadcast(text)
|
||
await update.message.reply_text(
|
||
f"📤 Рассылка завершена\n\n"
|
||
f"✅ Отправлено: {success}\n"
|
||
f"❌ Не доставлено: {failed}\n"
|
||
f"👥 Всего: {total}"
|
||
)
|
||
return
|
||
|
||
# Сохраняем chat_id админа при первом сообщении
|
||
chat_id = update.message.chat_id
|
||
saved_chat_id = get_admin_chat_id()
|
||
if saved_chat_id != chat_id:
|
||
save_admin_chat_id(chat_id)
|
||
if saved_chat_id is None:
|
||
await update.message.reply_text("✅ Админ бот активирован! Теперь вы будете получать все скачанные видео.\n\nДоступные команды: /stat, /send_messages")
|
||
else:
|
||
await update.message.reply_text("Это админский бот. Доступные команды: /stat, /send_messages")
|
||
else:
|
||
if update.message and update.message.text:
|
||
await update.message.reply_text("Это админский бот. Доступные команды: /stat, /send_messages")
|
||
|
||
|
||
def main():
|
||
"""Главная функция для запуска админ бота"""
|
||
if not ADMIN_BOT_TOKEN:
|
||
logger.error("ADMIN_BOT_TOKEN не установлен!")
|
||
return
|
||
|
||
# Создаем приложение
|
||
request = HTTPXRequest(
|
||
read_timeout=120,
|
||
connect_timeout=60
|
||
)
|
||
application = (
|
||
Application.builder()
|
||
.token(ADMIN_BOT_TOKEN)
|
||
.request(request)
|
||
.get_updates_request(HTTPXRequest(read_timeout=120, connect_timeout=60))
|
||
.build()
|
||
)
|
||
|
||
# Регистрируем обработчики
|
||
application.add_handler(CommandHandler("start", start_command))
|
||
application.add_handler(CommandHandler("stat", stat_command))
|
||
application.add_handler(CommandHandler("send_messages", send_messages_command))
|
||
application.add_handler(CommandHandler("cancel", cancel_command))
|
||
application.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, handle_message))
|
||
|
||
# Запускаем бота
|
||
logger.info("Админ бот запущен")
|
||
application.run_polling(allowed_updates=Update.ALL_TYPES)
|
||
|
||
|
||
if __name__ == '__main__':
|
||
main()
|
||
|