Compare commits

..

6 Commits

Author SHA1 Message Date
Yuriy Yuriev 19acea8481 fix 2026-07-21 00:22:37 +07:00
Yuriy Yuriev 9aa072792d fix 2026-07-21 00:05:25 +07:00
Yuriy Yuriev 4bd13ad041 fix 2026-07-18 18:43:27 +07:00
Yuriy Yuriev 8bdd0c784e fix 2026-07-18 18:33:54 +07:00
Yuriy Yuriev 6264b24f8f fix 2026-07-08 01:09:53 +07:00
Yuriy Yuriev 44ce29d2e0 fix 2026-07-05 16:59:33 +07:00
10 changed files with 508 additions and 304 deletions
+9 -1
View File
@@ -4,7 +4,15 @@
"Bash(pip install *)",
"Bash(python -c \"import psutil; print\\('psutil ok, version:', psutil.__version__\\)\")",
"Bash(python -c ' *)",
"WebFetch(domain:doc.heleket.com)"
"WebFetch(domain:doc.heleket.com)",
"Bash(chcp)",
"PowerShell(chcp)",
"PowerShell(python -c \"from core.logger import setup_logger; log = setup_logger\\('test'\\); log.info\\('?? @Nightbot : 15 '\\)\")",
"Read(//c/Users/user/AppData/Local/Programs/Python/**)",
"Bash(where.exe python *)",
"Bash(.venv/Scripts/python -c ' *)",
"Bash(/c/Work/PythonBot/.venv/Scripts/python.exe -c ' *)",
"PowerShell(& \"C:\\\\Work\\\\PythonBot\\\\.venv\\\\Scripts\\\\python.exe\" -c \"from core.logger import setup_logger; log = setup_logger\\('test'\\); log.info\\('test kirillicy klikov znak'\\)\")"
]
}
}
+15 -12
View File
@@ -1,8 +1,6 @@
import asyncio
import hashlib
import json
import secrets
import logging
import secrets
from datetime import datetime, timedelta
from pathlib import Path
from typing import Dict, Optional
@@ -16,7 +14,11 @@ _SESSIONS_FILE = Path("data/sessions.json")
class AuthManager:
"""Управляет авторизацией: хеширование паролей, сессии, верификация, роли."""
"""Управляет сессиями и ролями.
Личность пользователя подтверждает Telegram, поэтому пароля для входа нет.
ADMIN_PASSWORD остаётся единственным секретом — им повышают роль до admin.
"""
def __init__(self, session_timeout_minutes: int = 120):
self._sessions: Dict[int, datetime] = {}
@@ -28,15 +30,16 @@ class AuthManager:
self._restore_sessions()
@staticmethod
def _hash_password(password: str, salt: str = None) -> tuple:
if salt is None:
salt = secrets.token_hex(16)
result = hashlib.sha256((salt + password).encode()).hexdigest()
return result, salt
def verify_admin_password(password: str) -> bool:
"""Проверяет пароль администратора.
def verify_password(self, password: str, password_hash: str, salt: str) -> bool:
computed, _ = self._hash_password(password, salt)
return secrets.compare_digest(computed, password_hash)
Пустой ADMIN_PASSWORD означает, что админ-режим отключён — иначе
ненастроенный бот пускал бы в админку по пустой строке.
"""
if not ADMIN_PASSWORD:
logger.warning("ADMIN_PASSWORD not set — admin mode is disabled")
return False
return secrets.compare_digest(password, ADMIN_PASSWORD)
def is_authenticated(self, user_id: int) -> bool:
if user_id not in self._sessions:
+54 -25
View File
@@ -1,50 +1,74 @@
"""
Хранилище учётных данных пользователей.
Хранилище зарегистрированных пользователей.
"""
import asyncio
import json
import logging
from datetime import datetime
from pathlib import Path
from typing import Optional, Dict
logger = logging.getLogger(__name__)
class AuthStorage:
"""Хранилище паролей пользователей в JSON-файле."""
"""Реестр пользователей бота в JSON-файле.
Паролей не хранит: личность подтверждает Telegram, идентификатором
служит user_id.
"""
def __init__(self, file_path: str = "data/users.json"):
self.file_path = Path(file_path)
self.file_path.parent.mkdir(parents=True, exist_ok=True)
self._lock = asyncio.Lock()
async def register(self, user_id: int, password_hash: str, salt: str) -> bool:
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
if str(user_id) in data:
return False
data[str(user_id)] = {"password_hash": password_hash, "salt": salt}
await asyncio.to_thread(self._save_sync, data)
logger.info(f"User {user_id} registered")
return True
async def register(
self, user_id: int, username: str = None, full_name: str = None
) -> bool:
"""Регистрирует пользователя и освежает его профиль.
async def get_user(self, user_id: int) -> Optional[Dict]:
Профиль перезаписывается и для уже известных пользователей — ник в
Telegram может смениться, а взять его неоткуда, кроме входящего
апдейта.
"""
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
return data.get(str(user_id))
async def change_password(self, user_id: int, new_hash: str, new_salt: str) -> bool:
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
user_str = str(user_id)
if user_str not in data:
return False
data[user_str]["password_hash"] = new_hash
data[user_str]["salt"] = new_salt
key = str(user_id)
is_new = key not in data
record = data.get(key, {})
if is_new:
record["registered_at"] = datetime.now().isoformat()
if username is not None:
record["username"] = username
if full_name is not None:
record["full_name"] = full_name
data[key] = record
await asyncio.to_thread(self._save_sync, data)
logger.info(f"User {user_id} changed password")
return True
if is_new:
logger.info(f"User {user_id} registered ({username or full_name or 'no name'})")
return is_new
async def get_profile(self, user_id: int) -> dict:
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
return data.get(str(user_id), {})
async def get_all_profiles(self) -> dict:
"""Возвращает {user_id: профиль} для всех пользователей."""
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
return {int(k): v for k, v in data.items()}
@staticmethod
def display_name(profile: dict) -> str:
"""Человекочитаемое имя: @ник, иначе имя, иначе пусто."""
if not profile:
return ""
username = profile.get("username")
if username:
return f"@{username}"
return profile.get("full_name") or ""
async def delete_user(self, user_id: int) -> bool:
async with self._lock:
@@ -62,6 +86,11 @@ class AuthStorage:
data = await asyncio.to_thread(self._load_sync)
return str(user_id) in data
async def get_all_user_ids(self) -> list[int]:
async with self._lock:
data = await asyncio.to_thread(self._load_sync)
return [int(k) for k in data.keys()]
def _load_sync(self) -> dict:
if not self.file_path.exists():
return {}
+3
View File
@@ -28,6 +28,9 @@ def setup_logger(
datefmt='%Y-%m-%d %H:%M:%S'
)
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setLevel(log_level)
console_handler.setFormatter(console_formatter)
+313 -231
View File
@@ -33,7 +33,6 @@ from services.payment_service import HelketPayment
from services.twitch_api import get_viewer_count
from auth.manager import AuthManager
from auth.storage import AuthStorage
from auth import ADMIN_PASSWORD
logger = logging.getLogger(__name__)
@@ -69,12 +68,13 @@ class BotInterface:
self._topup_confirm: Dict[int, int] = {} # user_id -> rub amount pending confirm
self._user_task_state: Dict[int, dict] = {} # шаги создания задачи пользователем
self._admin_balance_state: Dict[int, int] = {} # admin_id -> target_user_id
self._admin_rub_state: Dict[int, int] = {} # admin_id -> target_user_id
self._admin_msg_state: Dict[int, int] = {} # admin_id -> target_user_id
self._menu_msg: Dict[int, int] = {} # user_id -> inline keyboard message_id
self._menu_top_msg: Dict[int, int] = {} # user_id -> reply keyboard message_id
self._chat_history: Dict[int, list] = {}
self._bot_ref = bot_ref
self._waiting_password: Dict[int, str] = {} # user_id -> "register" or "login"
self._waiting_password: Dict[int, str] = {} # user_id -> "admin_upgrade"
self._start_time = datetime.now()
async def _update_commands(self):
@@ -97,78 +97,18 @@ class BotInterface:
latest_msg_id=message.message_id
)
# Пароля для входа нет — личность подтверждает сам Telegram.
# Первый /start заодно регистрирует пользователя; для остальных
# обновляем ник, он мог смениться.
await interface._remember_profile(message.from_user)
if interface.auth_manager.is_authenticated(user_id) and interface.auth_manager.is_admin(user_id):
await interface._show_main_menu(message)
return
if interface.auth_manager.is_authenticated(user_id):
await interface._show_user_menu(message)
return
exists = await interface.auth_storage.user_exists(user_id)
if exists:
sent = await message.answer("С возвращением! Введите пароль:")
interface._waiting_password[user_id] = "login"
else:
sent = await message.answer(
"Добро пожаловать!\n\nПридумайте пароль для входа (мин. 4 символа):"
)
interface._waiting_password[user_id] = "register"
interface._menu_msg[user_id] = sent.message_id
@dp.message(Command("register"))
async def cmd_register(message: Message):
"""Регистрация: /register <пароль>"""
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
parts = message.text.split(maxsplit=1)
if len(parts) < 2 or len(parts[1]) < 4:
await interface._send_temp(message, "❌ Использование: /register <пароль от 4 символов>")
return
password = parts[1]
user_id = message.from_user.id
exists = await interface.auth_storage.user_exists(user_id)
if exists:
await interface._send_temp(message, "⚠️ Вы уже зарегистрированы. Используйте /login <пароль>")
return
password_hash, salt = interface.auth_manager._hash_password(password)
await interface.auth_storage.register(user_id, password_hash, salt)
is_admin = (password == ADMIN_PASSWORD)
interface.auth_manager.login(user_id, is_admin=is_admin)
if is_admin:
await interface._update_commands()
await interface._show_main_menu(message)
else:
await interface._show_user_menu(message)
@dp.message(Command("login"))
async def cmd_login(message: Message):
"""Вход: /login <пароль>"""
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
parts = message.text.split(maxsplit=1)
if len(parts) < 2:
await interface._send_temp(message, "❌ Использование: /login <пароль>")
return
password = parts[1]
user_id = message.from_user.id
user_data = await interface.auth_storage.get_user(user_id)
if not user_data:
await interface._send_temp(message, "❌ Вы не зарегистрированы.\nИспользуйте /register <пароль>")
return
if interface.auth_manager.is_locked_out(user_id):
await interface._send_temp(message, "🔒 Слишком много попыток. Попробуйте позже.")
return
if password == ADMIN_PASSWORD:
interface.auth_manager.login(user_id, is_admin=True)
await interface._update_commands()
await interface._show_main_menu(message)
elif interface.auth_manager.verify_password(
password, user_data["password_hash"], user_data["salt"]
):
if not interface.auth_manager.is_authenticated(user_id):
interface.auth_manager.login(user_id, is_admin=False)
await interface._show_user_menu(message)
else:
interface.auth_manager.record_failed_attempt(user_id)
await interface._send_temp(message, "❌ Неверный пароль")
@dp.message(Command("logout"))
async def cmd_logout(message: Message):
@@ -974,13 +914,6 @@ class BotInterface:
await callback.answer("❌ Недостаточно баланса", show_alert=True)
return
all_tasks = await interface.task_manager.get_all_tasks()
active_count = sum(
1 for p in all_tasks.values()
if p.user_id == uid and not p.completed and not p.stopped
)
if active_count >= 3:
await callback.answer("❌ Максимум 3 активные задачи", show_alert=True)
return
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="📺 Twitch мониторинг", callback_data="task_type_twitch"))
builder.row(InlineKeyboardButton(text="🔗 Посещение по ссылке", callback_data="task_type_url"))
@@ -1209,6 +1142,37 @@ class BotInterface:
if not await require_admin(callback):
return
await interface._show_users_list(callback.message, edit=True)
# === Листание админских списков ===
def _page_arg(data: str, prefix: str) -> int:
try:
return int(data.replace(prefix, "", 1))
except ValueError:
return 0
@dp.callback_query(F.data.startswith("users_page_"))
async def cb_users_page(callback: CallbackQuery):
if not await require_admin(callback):
return
page = _page_arg(callback.data, "users_page_")
await interface._show_users_list(callback.message, edit=True, page=page)
await callback.answer()
@dp.callback_query(F.data.startswith("streamers_page_"))
async def cb_streamers_page(callback: CallbackQuery):
if not await require_admin(callback):
return
page = _page_arg(callback.data, "streamers_page_")
await interface._show_streamers_list(callback.message, edit=True, page=page)
await callback.answer()
@dp.callback_query(F.data.startswith("tasks_page_"))
async def cb_tasks_page(callback: CallbackQuery):
if not await require_admin(callback):
return
page = _page_arg(callback.data, "tasks_page_")
await interface._show_tasks_list(callback.message, edit=True, page=page)
await callback.answer()
await callback.answer()
@dp.callback_query(F.data.startswith("udetail_"))
@@ -1241,6 +1205,24 @@ class BotInterface:
pass
await callback.answer()
@dp.callback_query(F.data.startswith("ubal_menu_"))
async def cb_user_balance_menu(callback: CallbackQuery):
if not await require_admin(callback):
return
target_uid = int(callback.data.replace("ubal_menu_", "", 1))
builder = InlineKeyboardBuilder()
for amount in [10, 50, 100, 500]:
builder.button(text=f"+{amount}", callback_data=f"ubal_{amount}_{target_uid}")
builder.adjust(4)
builder.row(InlineKeyboardButton(text="✏️ Другая сумма", callback_data=f"ubal_custom_{target_uid}"))
builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data=f"udetail_{target_uid}"))
clicks = await interface.balance_storage.get_balance(target_uid)
await callback.message.edit_text(
f"👤 {target_uid} — переходов: {clicks}\n\nДобавить переходов:",
reply_markup=builder.as_markup()
)
await callback.answer()
@dp.callback_query(F.data.startswith("ubal_"))
async def cb_user_balance(callback: CallbackQuery):
if not await require_admin(callback):
@@ -1252,7 +1234,7 @@ class BotInterface:
interface._admin_balance_state[callback.from_user.id] = target_uid
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data=f"udetail_{target_uid}"))
await callback.message.answer("Введите сумму для пополнения:", reply_markup=builder.as_markup())
await callback.message.answer("Введите количество переходов:", reply_markup=builder.as_markup())
await callback.answer()
else:
amount = int(parts[1])
@@ -1261,6 +1243,44 @@ class BotInterface:
await callback.answer(f"✅ +{amount} переходов. Баланс: {new_balance}")
await interface._show_user_detail(callback, target_uid)
@dp.callback_query(F.data.startswith("urub_menu_"))
async def cb_user_rub_menu(callback: CallbackQuery):
if not await require_admin(callback):
return
target_uid = int(callback.data.replace("urub_menu_", "", 1))
builder = InlineKeyboardBuilder()
for amount in [100, 500, 1000, 3000]:
builder.button(text=f"+{amount}", callback_data=f"urub_{amount}_{target_uid}")
builder.adjust(4)
builder.row(InlineKeyboardButton(text="✏️ Другая сумма", callback_data=f"urub_custom_{target_uid}"))
builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data=f"udetail_{target_uid}"))
rub = await interface.rub_storage.get_balance(target_uid)
await callback.message.edit_text(
f"👤 {target_uid} — рублей: {rub}\n\nДобавить рублей:",
reply_markup=builder.as_markup()
)
await callback.answer()
@dp.callback_query(F.data.startswith("urub_"))
async def cb_user_rub(callback: CallbackQuery):
if not await require_admin(callback):
return
parts = callback.data.split("_")
# формат: urub_{amount}_{user_id} или urub_custom_{user_id}
if parts[1] == "custom":
target_uid = int(parts[2])
interface._admin_rub_state[callback.from_user.id] = target_uid
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data=f"udetail_{target_uid}"))
await callback.message.answer("Введите сумму в рублях:", reply_markup=builder.as_markup())
await callback.answer()
else:
amount = int(parts[1])
target_uid = int(parts[2])
new_balance = await interface.rub_storage.add_balance(target_uid, amount)
await callback.answer(f"✅ +{amount} ₽. Баланс: {new_balance}")
await interface._show_user_detail(callback, target_uid)
@dp.callback_query(F.data.startswith("utasks_"))
async def cb_user_tasks_admin(callback: CallbackQuery):
if not await require_admin(callback):
@@ -1295,59 +1315,24 @@ class BotInterface:
await message.answer(f"❌ Не удалось отправить: {e}")
return
# Ввод пароля (неавторизованный пользователь ждёт пароль)
# Ввод пароля администратора (единственный оставшийся пароль)
if user_id in interface._waiting_password:
if not text:
return
# Пароль удаляем немедленно
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
mode = interface._waiting_password.pop(user_id)
interface._waiting_password.pop(user_id)
if mode == "register":
if len(text) < 4:
sent = await message.answer("Пароль слишком короткий, мин. 4 символа:")
interface._menu_msg[user_id] = sent.message_id
interface._waiting_password[user_id] = "register"
return
password_hash, salt = interface.auth_manager._hash_password(text)
await interface.auth_storage.register(user_id, password_hash, salt)
is_admin = (text == ADMIN_PASSWORD)
interface.auth_manager.login(user_id, is_admin=is_admin)
if is_admin:
await interface._update_commands()
await interface._show_main_menu(message)
else:
await interface._show_user_menu(message)
elif mode == "login":
user_data = await interface.auth_storage.get_user(user_id)
if not user_data:
await interface._send_temp(message, "Вы не зарегистрированы")
return
if interface.auth_manager.is_locked_out(user_id):
await interface._send_temp(message, "Слишком много попыток, попробуйте позже")
await interface._send_temp(message, "🔒 Слишком много попыток. Попробуйте позже.")
return
if text == ADMIN_PASSWORD:
if interface.auth_manager.verify_admin_password(text):
interface.auth_manager.login(user_id, is_admin=True)
await interface._update_commands()
await interface._show_main_menu(message)
elif interface.auth_manager.verify_password(
text, user_data["password_hash"], user_data["salt"]
):
interface.auth_manager.login(user_id, is_admin=False)
await interface._show_user_menu(message)
else:
interface.auth_manager.record_failed_attempt(user_id)
sent = await message.answer("Неверный пароль, попробуйте ещё раз:")
interface._menu_msg[user_id] = sent.message_id
interface._waiting_password[user_id] = "login"
elif mode == "admin_upgrade":
if text == ADMIN_PASSWORD:
interface.auth_manager.login(user_id, is_admin=True)
await interface._update_commands()
await interface._show_main_menu(message)
else:
sent = await message.answer("Неверный пароль администратора:")
interface._menu_msg[user_id] = sent.message_id
interface._waiting_password[user_id] = "admin_upgrade"
@@ -1376,7 +1361,7 @@ class BotInterface:
await interface._show_user_menu(message)
return
# Ввод суммы баланса (admin)
# Ввод количества переходов (admin)
if user_id in interface._admin_balance_state:
target_uid = interface._admin_balance_state.pop(user_id)
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
@@ -1385,7 +1370,21 @@ class BotInterface:
if amount <= 0:
raise ValueError
new_balance = await interface.balance_storage.add_balance(target_uid, amount)
await interface._send_temp(message, f"Пополнено на {amount}. Баланс: {new_balance}")
await interface._send_temp(message, f"+{amount} переходов. Баланс: {new_balance}")
except (ValueError, TypeError):
await interface._send_temp(message, "❌ Введите целое положительное число")
return
# Ввод рублёвого пополнения (admin)
if user_id in interface._admin_rub_state:
target_uid = interface._admin_rub_state.pop(user_id)
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
try:
amount = int(text)
if amount <= 0:
raise ValueError
new_balance = await interface.rub_storage.add_balance(target_uid, amount)
await interface._send_temp(message, f"✅ +{amount} ₽. Баланс: {new_balance}")
except (ValueError, TypeError):
await interface._send_temp(message, "❌ Введите целое положительное число")
return
@@ -1596,6 +1595,45 @@ class BotInterface:
await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id))
await self._safe_delete(message.bot, message.chat.id, message.message_id)
# Записей на страницу. Telegram режет сообщение на 4096 символах и
# 100 кнопках — поэтому списки листаем, а не показываем целиком.
_PAGE_SIZE = 8
@staticmethod
def _paginate(items: list, page: int, page_size: int = None) -> tuple:
"""Возвращает (срез страницы, нормализованный номер, всего страниц)."""
page_size = page_size or BotInterface._PAGE_SIZE
total_pages = max(1, (len(items) + page_size - 1) // page_size)
page = max(0, min(page, total_pages - 1))
start = page * page_size
return items[start:start + page_size], page, total_pages
@staticmethod
def _add_nav(builder, prefix: str, page: int, total_pages: int) -> None:
"""Добавляет строку навигации «◀ 2/5 ▶», если страниц больше одной."""
if total_pages <= 1:
return
row = []
if page > 0:
row.append(InlineKeyboardButton(text="◀️", callback_data=f"{prefix}{page - 1}"))
row.append(InlineKeyboardButton(text=f"{page + 1}/{total_pages}", callback_data="noop"))
if page < total_pages - 1:
row.append(InlineKeyboardButton(text="▶️", callback_data=f"{prefix}{page + 1}"))
builder.row(*row)
async def _remember_profile(self, from_user) -> None:
"""Сохраняет/освежает ник пользователя из входящего апдейта."""
if not from_user:
return
try:
await self.auth_storage.register(
from_user.id,
username=from_user.username or "",
full_name=(from_user.full_name or "").strip(),
)
except Exception as e:
logger.warning(f"Failed to remember profile for {from_user.id}: {e}")
async def _edit_or_send(self, message: Message, text: str, markup, edit: bool, user_id: int = None) -> None:
if edit:
try:
@@ -1645,116 +1683,100 @@ class BotInterface:
await self._track_msg(uid, sent1.message_id)
await self._track_msg(uid, sent2.message_id)
async def _show_streamers_list(self, message: Message, edit: bool = False):
async def _show_streamers_list(self, message: Message, edit: bool = False, page: int = 0):
active = await self.task_manager.get_active_tasks()
completed = await self.task_manager.get_completed_tasks()
twitch_active = {tid: p for tid, p in active.items() if p.task_type == "twitch_irc"}
twitch_done = {tid: p for tid, p in completed.items() if p.task_type == "twitch_irc"}
# Плоский список секция→задача: листаем сквозь обе секции,
# заголовок печатаем при смене секции внутри страницы
entries = [("🔄 Активные", tid, p) for tid, p in active.items() if p.task_type == "twitch_irc"]
entries += [("📁 Завершённые", tid, p) for tid, p in completed.items() if p.task_type == "twitch_irc"]
if not twitch_active and not twitch_done:
builder = InlineKeyboardBuilder()
if not entries:
text = "📺 Стримеры\n\nЗадач нет. Отправьте название канала для добавления."
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
else:
text = "📺 Стримеры\n\n"
builder = InlineKeyboardBuilder()
page_items, page, total_pages = self._paginate(entries, page)
text = f"📺 Стримеры ({len(entries)})"
if total_pages > 1:
text += f" — страница {page + 1}/{total_pages}"
text += "\n\n"
if twitch_active:
text += "🔄 Активные:\n"
for tid, p in list(twitch_active.items())[:10]:
current_section = None
for section, tid, p in page_items:
if section != current_section:
text += f"{section}:\n"
current_section = section
emoji = p.get_status_emoji()
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все"
text += (
f"{emoji} {p.channel}\n"
f" 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n"
)
suffix = f"🔗{p.links_found}" if section.startswith("🔄") else "завершена"
builder.row(InlineKeyboardButton(
text=f"{emoji} {p.channel} (🔗{p.links_found})",
text=f"{emoji} {p.channel} ({suffix})",
callback_data=f"tdetail_{tid}"
))
if twitch_done:
text += "\n📁 Завершенные:\n"
for tid, p in list(twitch_done.items())[:5]:
emoji = p.get_status_emoji()
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все"
text += (
f"{emoji} {p.channel}\n"
f" 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n"
)
builder.row(InlineKeyboardButton(
text=f"{emoji} {p.channel} (завершена)",
callback_data=f"tdetail_{tid}"
))
self._add_nav(builder, "streamers_page_", page, total_pages)
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
await self._edit_or_send(message, text, builder.as_markup(), edit)
async def _show_tasks_list(self, message: Message, edit: bool = False):
async def _show_tasks_list(self, message: Message, edit: bool = False, page: int = 0):
"""Все задачи — посещения и мониторинг."""
active = await self.task_manager.get_active_tasks()
completed = await self.task_manager.get_completed_tasks()
visit_active = {tid: p for tid, p in active.items() if p.task_type in ("visit", "user_visit")}
visit_done = {tid: p for tid, p in completed.items() if p.task_type in ("visit", "user_visit")}
twitch_active = {tid: p for tid, p in active.items() if p.task_type == "twitch_irc"}
twitch_done = {tid: p for tid, p in completed.items() if p.task_type == "twitch_irc"}
def _pick(src, types):
return [(tid, p) for tid, p in src.items() if p.task_type in types]
# Плоский список во всех четырёх секциях — листаем сквозной пагинацией
entries = [("📺 Twitch — активные", t, p) for t, p in _pick(active, ("twitch_irc",))]
entries += [("📺 Twitch — завершённые", t, p) for t, p in _pick(completed, ("twitch_irc",))]
entries += [("🔗 Посещения — активные", t, p) for t, p in _pick(active, ("visit", "user_visit"))]
entries += [("🔗 Посещения — завершённые", t, p) for t, p in _pick(completed, ("visit", "user_visit"))]
builder = InlineKeyboardBuilder()
if not any([visit_active, visit_done, twitch_active, twitch_done]):
if not entries:
text = "📊 Задачи\n\nЗадач нет."
else:
text = "📊 Задачи\n\n"
page_items, page, total_pages = self._paginate(entries, page)
text = f"📊 Задачи ({len(entries)})"
if total_pages > 1:
text += f" — страница {page + 1}/{total_pages}"
text += "\n\n"
if twitch_active:
text += "📺 Twitch — активные:\n"
for tid, p in list(twitch_active.items())[:10]:
current_section = None
for section, tid, p in page_items:
if section != current_section:
text += f"\n{section}:\n" if current_section else f"{section}:\n"
current_section = section
em = p.get_status_emoji()
is_active = "активные" in section
if section.startswith("📺"):
if is_active:
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все"
text += f"{em} {p.channel}\n 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n"
builder.row(InlineKeyboardButton(
text=f"{em} 📺 {p.channel} (🔗{p.links_found})",
callback_data=f"tdetail_{tid}"
))
if twitch_done:
text += "\n📺 Twitch — завершённые:\n"
for tid, p in list(twitch_done.items())[:5]:
em = p.get_status_emoji()
label = f"{em} 📺 {p.channel} (🔗{p.links_found})"
else:
text += f"{em} {p.channel}\n"
builder.row(InlineKeyboardButton(
text=f"{em} 📺 {p.channel} (завершена)",
callback_data=f"tdetail_{tid}"
))
if visit_active:
text += "\n🔗 Посещения — активные:\n"
for tid, p in list(visit_active.items())[:10]:
em = p.get_status_emoji()
label = f"{em} 📺 {p.channel} (завершена)"
else:
url_s = (p.url or "")[:35]
uid_tag = f" | uid:{p.user_id}" if p.user_id else ""
if is_active:
text += f"{em} {url_s}{uid_tag}\n{p.successful_visits}/{p.max_visits} | 📊{p.total_visits}\n"
builder.row(InlineKeyboardButton(
text=f"{em} 🔗 {url_s[:28]} ({p.successful_visits}/{p.max_visits})",
callback_data=f"tdetail_{tid}"
))
if visit_done:
text += "\n🔗 Посещения — завершённые:\n"
for tid, p in list(visit_done.items())[:5]:
em = p.get_status_emoji()
url_s = (p.url or "")[:35]
uid_tag = f" | uid:{p.user_id}" if p.user_id else ""
label = f"{em} 🔗 {url_s[:28]} ({p.successful_visits}/{p.max_visits})"
else:
text += f"{em} {url_s}{uid_tag}\n"
builder.row(InlineKeyboardButton(
text=f"{em} 🔗 {url_s[:28]} (завершена)",
callback_data=f"tdetail_{tid}"
))
label = f"{em} 🔗 {url_s[:28]} (завершена)"
builder.row(InlineKeyboardButton(text=label, callback_data=f"tdetail_{tid}"))
self._add_nav(builder, "tasks_page_", page, total_pages)
builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_tasks"))
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
@@ -2125,8 +2147,29 @@ class BotInterface:
reply_markup=builder.as_markup()
)
async def _drop_pending_series(self, params: TaskParams, url_tasks: Optional[set]) -> int:
"""Сбрасывает очередь кликов, накопленную за прошлую трансляцию.
Отменяет незавершённые серии `_process_url` и обнуляет счётчик, чтобы
клики по устаревшим ссылкам не «догоняли» уже новый стрим.
Возвращает число отброшенных кликов.
"""
dropped = params.pending_visits
if url_tasks:
pending = [t for t in url_tasks if not t.done()]
for t in pending:
t.cancel()
if pending:
await asyncio.gather(*pending, return_exceptions=True)
# Отброшенные клики не будут выполнены — убираем их и из плана,
# иначе счётчик "выполнено/запланировано" никогда не сойдётся
params.total_planned_visits = max(0, params.total_planned_visits - dropped)
params.pending_visits = 0
params.remaining_clicks = 0
return dropped
async def _stream_watcher(self, task_id: str, params: TaskParams, message: Message,
check_interval: int = 60) -> None:
check_interval: int = 60, url_tasks: Optional[set] = None) -> None:
"""Проверяет онлайн-статус стрима, ставит/снимает паузу и сохраняет состояние."""
was_offline = params.stream_offline
@@ -2161,6 +2204,10 @@ class BotInterface:
elif viewers > 0 and was_offline:
was_offline = False
# Началась новая трансляция — чистим очередь прошлой.
# Делаем это до снятия паузы, иначе отложенные серии успеют
# проснуться и отработать по неактуальным ссылкам.
dropped = await self._drop_pending_series(params, url_tasks)
params.stream_offline = False
if params.auto_paused:
params.paused = False
@@ -2168,11 +2215,14 @@ class BotInterface:
await self.task_manager.update_task(task_id, paused=False)
if self.storage:
await self.storage.save_task(task_id, params)
logger.info(f"[watcher] {params.channel} is online ({viewers} viewers)")
msg = await send_message_safe(
message.bot, params.chat_id,
f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён"
logger.info(
f"[watcher] {params.channel} is online ({viewers} viewers), "
f"queue cleared: {dropped} clicks dropped"
)
text = f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён"
if dropped:
text += f"\n🧹 Очередь прошлой трансляции очищена: {dropped} кликов отброшено"
msg = await send_message_safe(message.bot, params.chat_id, text)
if msg and params.user_id:
await self._track_msg(params.user_id, msg.message_id)
@@ -2204,7 +2254,9 @@ class BotInterface:
logger.info(f"🚀 GQL chat monitor: {params.channel}")
watcher = asyncio.create_task(self._stream_watcher(task_id, params, message))
watcher = asyncio.create_task(
self._stream_watcher(task_id, params, message, url_tasks=url_tasks)
)
try:
await irc.listen_for_messages(
@@ -2361,6 +2413,7 @@ class BotInterface:
async def _process_url(self, url: str, username: str, task_id: str, params: TaskParams, message: Message):
"""Обрабатывает найденную ссылку (выполняется параллельно)."""
notif_msg = None
try:
params.links_found += 1
@@ -2389,7 +2442,7 @@ class BotInterface:
params.pending_visits += visits_count
params.total_planned_visits += visits_count
logger.info(f"🔗 @{username} {url}{visits_count} clicks ({calc_info}), total pending: {params.pending_visits}")
await send_message_safe(
notif_msg = await send_message_safe(
message.bot, message.chat.id,
f"🔗 @{username}: {url[:60]}\n"
f"📊 {calc_info}\n"
@@ -2402,6 +2455,7 @@ class BotInterface:
# Resume > 5 мин: серия этой ссылки полностью отбрасывается
SERIES_EXPIRY = 5 * 60
remaining = visits_count
successful_this_link = 0
paused_at: Optional[float] = None
first_click = True
@@ -2443,6 +2497,7 @@ class BotInterface:
return
if result.success:
params.successful_visits += 1
successful_this_link += 1
remaining -= 1
if params.user_id:
new_balance = await self.balance_storage.deduct(params.user_id, 1)
@@ -2470,9 +2525,23 @@ class BotInterface:
if params.stopped:
return
except asyncio.CancelledError:
# Очередь сброшена (началась новая трансляция) — серия прекращается
logger.info(f"🧹 Серия отменена: {url[:40]}")
raise
except Exception as e:
logger.error(f"Process URL error: {e}")
finally:
# Уведомление о серии убираем при любом исходе: завершение, отмена,
# обрыв прокси, исчерпание баланса
if notif_msg:
try:
await message.bot.delete_message(
chat_id=message.chat.id,
message_id=notif_msg.message_id,
)
except Exception:
pass
# =========================================================================
# КАБИНЕТ ПОЛЬЗОВАТЕЛЯ
@@ -2647,18 +2716,6 @@ class BotInterface:
await message.answer("❌ Неверное название канала (только буквы, цифры, _). Попробуйте ещё раз:")
return
channel = text.lower()
# Проверка на дубликат
all_tasks = await self.task_manager.get_all_tasks()
duplicate = any(
p.user_id == user_id and p.channel == channel and not p.completed and not p.stopped
for p in all_tasks.values()
)
if duplicate:
await self._send_temp(
message,
f"❌ У вас уже есть активная задача для канала `{channel}`.\nВведите другой канал:"
)
return
# Удаляем промпт шага 1 перед отправкой шага 2
if user_id in self._menu_msg:
await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id))
@@ -3069,51 +3126,76 @@ class BotInterface:
else:
await self._send_temp(message, "Ошибка запуска задачи")
async def _show_users_list(self, message: Message, edit: bool = False):
balances = await self.balance_storage.get_all()
all_tasks = await self.task_manager.get_all_tasks()
async def _show_users_list(self, message: Message, edit: bool = False, page: int = 0):
profiles, all_tasks, all_clicks, all_rub = await asyncio.gather(
self.auth_storage.get_all_profiles(),
self.task_manager.get_all_tasks(),
self.balance_storage.get_all(),
self.rub_storage.get_all(),
)
all_user_ids = list(profiles.keys())
text = "👥 Пользователи\n\n"
builder = InlineKeyboardBuilder()
if not balances:
text += "Нет пользователей с балансом."
if not all_user_ids:
text = "👥 Пользователи (0)\n\nНет зарегистрированных пользователей."
else:
for uid_str, balance in list(balances.items())[:15]:
uid = int(uid_str)
page_ids, page, total_pages = self._paginate(all_user_ids, page)
text = f"👥 Пользователи ({len(all_user_ids)})"
if total_pages > 1:
text += f" — страница {page + 1}/{total_pages}"
text += "\n\n"
for uid in page_ids:
clicks = all_clicks.get(str(uid), 0)
rub = all_rub.get(str(uid), 0)
active = sum(
1 for p in all_tasks.values()
if p.user_id == uid and not p.completed and not p.stopped
)
text += f"👤 `{uid}` — 💰 {balance} | 📋 {active} задач\n"
bal_str = f"🖱{clicks}"
if rub:
bal_str += f" 🪙{rub}"
if active:
bal_str += f" 📋{active}"
# Ник известен только после того, как пользователь напишет боту
name = AuthStorage.display_name(profiles.get(uid))
text += f"👤 {name + ' ' if name else ''}`{uid}` — {bal_str}\n"
builder.row(InlineKeyboardButton(
text=f"👤 {uid} (💰{balance})",
text=f"👤 {name or uid} {bal_str}",
callback_data=f"udetail_{uid}",
))
self._add_nav(builder, "users_page_", page, total_pages)
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
await self._edit_or_send(message, text, builder.as_markup(), edit)
async def _show_user_detail(self, callback: CallbackQuery, target_user_id: int):
balance = await self.balance_storage.get_balance(target_user_id)
clicks = await self.balance_storage.get_balance(target_user_id)
rub = await self.rub_storage.get_balance(target_user_id)
profile = await self.auth_storage.get_profile(target_user_id)
all_tasks = await self.task_manager.get_all_tasks()
user_tasks = {tid: p for tid, p in all_tasks.items() if p.user_id == target_user_id}
active = sum(1 for p in user_tasks.values() if not p.completed and not p.stopped)
name = AuthStorage.display_name(profile)
header = f"👤 {name} ({target_user_id})" if name else f"👤 Пользователь {target_user_id}"
full_name = profile.get("full_name")
if full_name and name != full_name:
header += f"\n📝 {full_name}"
text = (
f"👤 Пользователь {target_user_id}\n\n"
f"🔢 Переходов: {balance}\n"
f"📋 Задач: {len(user_tasks)} (активных: {active})\n\n"
"Пополнить баланс:"
f"{header}\n\n"
f"🖱 Переходов: {clicks}\n"
f"🪙 Рублей: {rub}\n"
f"📋 Задач: {len(user_tasks)} (активных: {active})\n"
)
builder = InlineKeyboardBuilder()
for amount in [10, 50, 100, 500]:
builder.button(text=f"+{amount}", callback_data=f"ubal_{amount}_{target_user_id}")
builder.adjust(4)
builder.row(InlineKeyboardButton(
text="✏️ Другая сумма",
callback_data=f"ubal_custom_{target_user_id}",
))
builder.row(InlineKeyboardButton(text=" Переходы", callback_data=f"ubal_menu_{target_user_id}"))
builder.row(InlineKeyboardButton(text=" Рубли", callback_data=f"urub_menu_{target_user_id}"))
if user_tasks:
builder.row(InlineKeyboardButton(
text=f"📋 Задачи ({len(user_tasks)})",
+42 -2
View File
@@ -6,7 +6,7 @@ import asyncio
import json
import logging
from pathlib import Path
from datetime import datetime
from datetime import datetime, timedelta
from typing import Dict, Optional
from managers.task_manager import TaskParams
@@ -309,6 +309,9 @@ class RubleBalanceStorage(_IntBalanceStorage):
class PaymentStorage:
"""Хранилище ожидающих платежей."""
# Сколько держать обработанные платежи ради защиты от дублей вебхука
PROCESSED_TTL_DAYS = 7
def __init__(self, file_path: str = "data/payments.json"):
self.file_path = Path(file_path)
self.file_path.parent.mkdir(parents=True, exist_ok=True)
@@ -327,9 +330,29 @@ class PaymentStorage:
with open(self.file_path, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
def _prune_processed(self, payments: dict) -> dict:
"""Удаляет обработанные платежи старше PROCESSED_TTL_DAYS.
Без этого файл рос бы бесконечно, а его читают при каждой
отрисовке кабинета.
"""
cutoff = datetime.now() - timedelta(days=self.PROCESSED_TTL_DAYS)
kept = {}
for pid, data in payments.items():
processed_at = data.get("processed_at")
if data.get("processed") and processed_at:
try:
if datetime.fromisoformat(processed_at) < cutoff:
continue
except ValueError:
pass # некорректная дата — запись оставляем
kept[pid] = data
return kept
async def save(self, payment_id: str, data: dict) -> None:
async with self._lock:
payments = await asyncio.to_thread(self._load_sync)
payments = self._prune_processed(payments)
payments[payment_id] = data
await asyncio.to_thread(self._save_sync, payments)
@@ -344,11 +367,28 @@ class PaymentStorage:
payments.pop(payment_id, None)
await asyncio.to_thread(self._save_sync, payments)
async def pop(self, payment_id: str) -> Optional[dict]:
"""Get and mark payment as processed. Returns None if not found or already processed."""
async with self._lock:
payments = await asyncio.to_thread(self._load_sync)
data = payments.get(payment_id)
if data is None or data.get("processed"):
return None
payments[payment_id]["processed"] = True
payments[payment_id]["processed_at"] = datetime.now().isoformat(timespec="seconds")
await asyncio.to_thread(self._save_sync, payments)
return data
async def get_by_user(self, user_id: int) -> Optional[dict]:
"""Возвращает незакрытый счёт пользователя.
Обработанные платежи остаются в файле как защита от повторного
начисления по дублю вебхука, но открытым счётом уже не считаются.
"""
async with self._lock:
payments = await asyncio.to_thread(self._load_sync)
for pid, data in payments.items():
if data.get("user_id") == user_id:
if data.get("user_id") == user_id and not data.get("processed"):
return {**data, "payment_id": pid}
return None
+1 -1
View File
@@ -332,7 +332,7 @@ class BrowserService:
current = new_url
redirect = True
task = asyncio.create_task(
page.wait_for_load_state("load", timeout=30000)
page.wait_for_load_state("domcontentloaded", timeout=30000)
)
await self._move_mouse(page, task)
await task
+11
View File
@@ -5,6 +5,7 @@ Twitch чат через IRC over WebSocket (wss://irc-ws.chat.twitch.tv:443).
import asyncio
import logging
import random
import re
from typing import Optional, Callable, Awaitable, List
from datetime import datetime
@@ -34,7 +35,13 @@ class TwitchGQLChatClient:
):
self.channel = channel.lower()
self.target_username = target_username.lower()
# Anonymous Twitch login (justinfanNNNNN) is shared across all clients by
# default, so parallel monitors on the same nick get kicked by Twitch
# (one connection per nick) — randomize it unless a real login is configured.
if irc_username or settings.IRC_USERNAME != "justinfan12345":
self.irc_username = irc_username or settings.IRC_USERNAME
else:
self.irc_username = f"justinfan{random.randint(10000, 99999)}"
self.irc_oauth = irc_oauth or settings.IRC_OAUTH
self._ws: Optional[aiohttp.ClientWebSocketResponse] = None
@@ -257,6 +264,10 @@ class TwitchGQLChatClient:
raise
except Exception:
pass
else:
logger.info(
f"🚫 Домен не разрешён: {url} (allowed_domains={allowed_domains})"
)
return stats
+7
View File
@@ -3,6 +3,7 @@ IRC сервис для мониторинга Twitch чата.
"""
import asyncio
import random
import re
import logging
from typing import Optional, Callable, Awaitable, List
@@ -28,7 +29,13 @@ class TwitchIRCClient:
):
self.channel = channel.lower()
self.target_username = target_username.lower()
# Anonymous Twitch login (justinfanNNNNN) is shared across all clients by
# default, so parallel monitors on the same nick get kicked by Twitch
# (one connection per nick) — randomize it unless a real login is configured.
if irc_username or settings.IRC_USERNAME != "justinfan12345":
self.irc_username = irc_username or settings.IRC_USERNAME
else:
self.irc_username = f"justinfan{random.randint(10000, 99999)}"
self.irc_oauth = irc_oauth or settings.IRC_OAUTH
self._reader: Optional[asyncio.StreamReader] = None
+24 -3
View File
@@ -11,6 +11,7 @@ import hashlib
import json
import logging
from datetime import datetime
from pathlib import Path
from typing import Optional
from aiohttp import web
@@ -74,6 +75,21 @@ class HelketWebhookServer:
await self._runner.cleanup()
logger.info("Heleket webhook server stopped")
async def _save_lost_webhook(self, payload: dict) -> None:
"""Сохраняет необработанный вебхук в файл и уведомляет админа."""
lost_file = Path("data/lost_webhooks.json")
try:
data = []
if lost_file.exists():
with open(lost_file, "r", encoding="utf-8") as f:
data = json.load(f)
data.append({**payload, "_saved_at": datetime.now().isoformat(timespec="seconds")})
with open(lost_file, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
logger.info(f"Lost webhook saved: uuid={payload.get('uuid')}")
except Exception as e:
logger.error(f"Failed to save lost webhook: {e}")
async def _handle_health(self, request: web.Request) -> web.Response:
return web.Response(text="ok")
@@ -102,9 +118,15 @@ class HelketWebhookServer:
logger.info(f"Webhook: uuid={payment_uuid} status={status} — ignored")
return web.Response(text="ok")
payment = await self._payment_storage.get(payment_uuid)
existing = await self._payment_storage.get(payment_uuid)
if existing and existing.get("processed"):
logger.info(f"Webhook: payment {payment_uuid} already processed — skipping duplicate")
return web.Response(text="ok")
payment = await self._payment_storage.pop(payment_uuid)
if not payment:
logger.warning(f"Webhook: payment {payment_uuid} not found (already processed?)")
logger.warning(f"Webhook: payment {payment_uuid} not found in storage")
await self._save_lost_webhook(payload)
return web.Response(text="ok")
user_id = payment.get("user_id")
@@ -147,7 +169,6 @@ class HelketWebhookServer:
}
logger.info(f"Webhook: user={user_id} +{rub} RUB")
await self._payment_storage.delete(payment_uuid)
await self._history_storage.add(user_id, history_record)
if chat_id: