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(pip install *)",
"Bash(python -c \"import psutil; print\\('psutil ok, version:', psutil.__version__\\)\")", "Bash(python -c \"import psutil; print\\('psutil ok, version:', psutil.__version__\\)\")",
"Bash(python -c ' *)", "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 json
import secrets
import logging import logging
import secrets
from datetime import datetime, timedelta from datetime import datetime, timedelta
from pathlib import Path from pathlib import Path
from typing import Dict, Optional from typing import Dict, Optional
@@ -16,7 +14,11 @@ _SESSIONS_FILE = Path("data/sessions.json")
class AuthManager: class AuthManager:
"""Управляет авторизацией: хеширование паролей, сессии, верификация, роли.""" """Управляет сессиями и ролями.
Личность пользователя подтверждает Telegram, поэтому пароля для входа нет.
ADMIN_PASSWORD остаётся единственным секретом — им повышают роль до admin.
"""
def __init__(self, session_timeout_minutes: int = 120): def __init__(self, session_timeout_minutes: int = 120):
self._sessions: Dict[int, datetime] = {} self._sessions: Dict[int, datetime] = {}
@@ -28,15 +30,16 @@ class AuthManager:
self._restore_sessions() self._restore_sessions()
@staticmethod @staticmethod
def _hash_password(password: str, salt: str = None) -> tuple: def verify_admin_password(password: str) -> bool:
if salt is None: """Проверяет пароль администратора.
salt = secrets.token_hex(16)
result = hashlib.sha256((salt + password).encode()).hexdigest()
return result, salt
def verify_password(self, password: str, password_hash: str, salt: str) -> bool: Пустой ADMIN_PASSWORD означает, что админ-режим отключён — иначе
computed, _ = self._hash_password(password, salt) ненастроенный бот пускал бы в админку по пустой строке.
return secrets.compare_digest(computed, password_hash) """
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: def is_authenticated(self, user_id: int) -> bool:
if user_id not in self._sessions: if user_id not in self._sessions:
+54 -25
View File
@@ -1,50 +1,74 @@
""" """
Хранилище учётных данных пользователей. Хранилище зарегистрированных пользователей.
""" """
import asyncio import asyncio
import json import json
import logging import logging
from datetime import datetime
from pathlib import Path from pathlib import Path
from typing import Optional, Dict
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
class AuthStorage: class AuthStorage:
"""Хранилище паролей пользователей в JSON-файле.""" """Реестр пользователей бота в JSON-файле.
Паролей не хранит: личность подтверждает Telegram, идентификатором
служит user_id.
"""
def __init__(self, file_path: str = "data/users.json"): def __init__(self, file_path: str = "data/users.json"):
self.file_path = Path(file_path) self.file_path = Path(file_path)
self.file_path.parent.mkdir(parents=True, exist_ok=True) self.file_path.parent.mkdir(parents=True, exist_ok=True)
self._lock = asyncio.Lock() self._lock = asyncio.Lock()
async def register(self, user_id: int, password_hash: str, salt: str) -> bool: async def register(
async with self._lock: self, user_id: int, username: str = None, full_name: str = None
data = await asyncio.to_thread(self._load_sync) ) -> bool:
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 get_user(self, user_id: int) -> Optional[Dict]: Профиль перезаписывается и для уже известных пользователей — ник в
Telegram может смениться, а взять его неоткуда, кроме входящего
апдейта.
"""
async with self._lock: async with self._lock:
data = await asyncio.to_thread(self._load_sync) data = await asyncio.to_thread(self._load_sync)
return data.get(str(user_id)) key = str(user_id)
is_new = key not in data
async def change_password(self, user_id: int, new_hash: str, new_salt: str) -> bool: record = data.get(key, {})
async with self._lock: if is_new:
data = await asyncio.to_thread(self._load_sync) record["registered_at"] = datetime.now().isoformat()
user_str = str(user_id) if username is not None:
if user_str not in data: record["username"] = username
return False if full_name is not None:
data[user_str]["password_hash"] = new_hash record["full_name"] = full_name
data[user_str]["salt"] = new_salt data[key] = record
await asyncio.to_thread(self._save_sync, data) await asyncio.to_thread(self._save_sync, data)
logger.info(f"User {user_id} changed password") if is_new:
return True 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 def delete_user(self, user_id: int) -> bool:
async with self._lock: async with self._lock:
@@ -62,6 +86,11 @@ class AuthStorage:
data = await asyncio.to_thread(self._load_sync) data = await asyncio.to_thread(self._load_sync)
return str(user_id) in data 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: def _load_sync(self) -> dict:
if not self.file_path.exists(): if not self.file_path.exists():
return {} return {}
+3
View File
@@ -28,6 +28,9 @@ def setup_logger(
datefmt='%Y-%m-%d %H:%M:%S' 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 = logging.StreamHandler(sys.stdout)
console_handler.setLevel(log_level) console_handler.setLevel(log_level)
console_handler.setFormatter(console_formatter) console_handler.setFormatter(console_formatter)
+340 -258
View File
@@ -33,7 +33,6 @@ from services.payment_service import HelketPayment
from services.twitch_api import get_viewer_count from services.twitch_api import get_viewer_count
from auth.manager import AuthManager from auth.manager import AuthManager
from auth.storage import AuthStorage from auth.storage import AuthStorage
from auth import ADMIN_PASSWORD
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -69,12 +68,13 @@ class BotInterface:
self._topup_confirm: Dict[int, int] = {} # user_id -> rub amount pending confirm self._topup_confirm: Dict[int, int] = {} # user_id -> rub amount pending confirm
self._user_task_state: Dict[int, dict] = {} # шаги создания задачи пользователем self._user_task_state: Dict[int, dict] = {} # шаги создания задачи пользователем
self._admin_balance_state: Dict[int, int] = {} # admin_id -> target_user_id 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._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_msg: Dict[int, int] = {} # user_id -> inline keyboard message_id
self._menu_top_msg: Dict[int, int] = {} # user_id -> reply keyboard message_id self._menu_top_msg: Dict[int, int] = {} # user_id -> reply keyboard message_id
self._chat_history: Dict[int, list] = {} self._chat_history: Dict[int, list] = {}
self._bot_ref = bot_ref 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() self._start_time = datetime.now()
async def _update_commands(self): async def _update_commands(self):
@@ -97,78 +97,18 @@ class BotInterface:
latest_msg_id=message.message_id 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): if interface.auth_manager.is_authenticated(user_id) and interface.auth_manager.is_admin(user_id):
await interface._show_main_menu(message) await interface._show_main_menu(message)
return return
if interface.auth_manager.is_authenticated(user_id): if not 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"]
):
interface.auth_manager.login(user_id, is_admin=False) interface.auth_manager.login(user_id, is_admin=False)
await interface._show_user_menu(message) await interface._show_user_menu(message)
else:
interface.auth_manager.record_failed_attempt(user_id)
await interface._send_temp(message, "❌ Неверный пароль")
@dp.message(Command("logout")) @dp.message(Command("logout"))
async def cmd_logout(message: Message): async def cmd_logout(message: Message):
@@ -973,14 +913,7 @@ class BotInterface:
if balance <= 0: if balance <= 0:
await callback.answer("❌ Недостаточно баланса", show_alert=True) await callback.answer("❌ Недостаточно баланса", show_alert=True)
return return
all_tasks = await interface.task_manager.get_all_tasks() 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 = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="📺 Twitch мониторинг", callback_data="task_type_twitch")) builder.row(InlineKeyboardButton(text="📺 Twitch мониторинг", callback_data="task_type_twitch"))
builder.row(InlineKeyboardButton(text="🔗 Посещение по ссылке", callback_data="task_type_url")) builder.row(InlineKeyboardButton(text="🔗 Посещение по ссылке", callback_data="task_type_url"))
@@ -1209,6 +1142,37 @@ class BotInterface:
if not await require_admin(callback): if not await require_admin(callback):
return return
await interface._show_users_list(callback.message, edit=True) 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() await callback.answer()
@dp.callback_query(F.data.startswith("udetail_")) @dp.callback_query(F.data.startswith("udetail_"))
@@ -1241,6 +1205,24 @@ class BotInterface:
pass pass
await callback.answer() 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_")) @dp.callback_query(F.data.startswith("ubal_"))
async def cb_user_balance(callback: CallbackQuery): async def cb_user_balance(callback: CallbackQuery):
if not await require_admin(callback): if not await require_admin(callback):
@@ -1252,7 +1234,7 @@ class BotInterface:
interface._admin_balance_state[callback.from_user.id] = target_uid interface._admin_balance_state[callback.from_user.id] = target_uid
builder = InlineKeyboardBuilder() builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data=f"udetail_{target_uid}")) 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() await callback.answer()
else: else:
amount = int(parts[1]) amount = int(parts[1])
@@ -1261,6 +1243,44 @@ class BotInterface:
await callback.answer(f"✅ +{amount} переходов. Баланс: {new_balance}") await callback.answer(f"✅ +{amount} переходов. Баланс: {new_balance}")
await interface._show_user_detail(callback, target_uid) 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_")) @dp.callback_query(F.data.startswith("utasks_"))
async def cb_user_tasks_admin(callback: CallbackQuery): async def cb_user_tasks_admin(callback: CallbackQuery):
if not await require_admin(callback): if not await require_admin(callback):
@@ -1295,62 +1315,27 @@ class BotInterface:
await message.answer(f"❌ Не удалось отправить: {e}") await message.answer(f"❌ Не удалось отправить: {e}")
return return
# Ввод пароля (неавторизованный пользователь ждёт пароль) # Ввод пароля администратора (единственный оставшийся пароль)
if user_id in interface._waiting_password: if user_id in interface._waiting_password:
if not text: if not text:
return return
# Пароль удаляем немедленно # Пароль удаляем немедленно
await interface._safe_delete(message.bot, message.chat.id, message.message_id) 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 interface.auth_manager.is_locked_out(user_id):
if len(text) < 4: await interface._send_temp(message, "🔒 Слишком много попыток. Попробуйте позже.")
sent = await message.answer("Пароль слишком короткий, мин. 4 символа:") return
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": if interface.auth_manager.verify_admin_password(text):
user_data = await interface.auth_storage.get_user(user_id) interface.auth_manager.login(user_id, is_admin=True)
if not user_data: await interface._update_commands()
await interface._send_temp(message, "Вы не зарегистрированы") await interface._show_main_menu(message)
return else:
if interface.auth_manager.is_locked_out(user_id): interface.auth_manager.record_failed_attempt(user_id)
await interface._send_temp(message, "Слишком много попыток, попробуйте позже") sent = await message.answer("Неверный пароль администратора:")
return interface._menu_msg[user_id] = sent.message_id
if text == ADMIN_PASSWORD: interface._waiting_password[user_id] = "admin_upgrade"
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"
return return
# Ввод шагов создания задачи (пользователь) # Ввод шагов создания задачи (пользователь)
@@ -1376,7 +1361,7 @@ class BotInterface:
await interface._show_user_menu(message) await interface._show_user_menu(message)
return return
# Ввод суммы баланса (admin) # Ввод количества переходов (admin)
if user_id in interface._admin_balance_state: if user_id in interface._admin_balance_state:
target_uid = interface._admin_balance_state.pop(user_id) target_uid = interface._admin_balance_state.pop(user_id)
await interface._safe_delete(message.bot, message.chat.id, message.message_id) await interface._safe_delete(message.bot, message.chat.id, message.message_id)
@@ -1385,7 +1370,21 @@ class BotInterface:
if amount <= 0: if amount <= 0:
raise ValueError raise ValueError
new_balance = await interface.balance_storage.add_balance(target_uid, amount) 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): except (ValueError, TypeError):
await interface._send_temp(message, "❌ Введите целое положительное число") await interface._send_temp(message, "❌ Введите целое положительное число")
return 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, self._menu_msg.pop(user_id))
await self._safe_delete(message.bot, message.chat.id, message.message_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: async def _edit_or_send(self, message: Message, text: str, markup, edit: bool, user_id: int = None) -> None:
if edit: if edit:
try: try:
@@ -1645,116 +1683,100 @@ class BotInterface:
await self._track_msg(uid, sent1.message_id) await self._track_msg(uid, sent1.message_id)
await self._track_msg(uid, sent2.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() active = await self.task_manager.get_active_tasks()
completed = await self.task_manager.get_completed_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Задач нет. Отправьте название канала для добавления." text = "📺 Стримеры\n\nЗадач нет. Отправьте название канала для добавления."
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
else: else:
text = "📺 Стримеры\n\n" page_items, page, total_pages = self._paginate(entries, page)
builder = InlineKeyboardBuilder() text = f"📺 Стримеры ({len(entries)})"
if total_pages > 1:
text += f" — страница {page + 1}/{total_pages}"
text += "\n\n"
if twitch_active: current_section = None
text += "🔄 Активные:\n" for section, tid, p in page_items:
for tid, p in list(twitch_active.items())[:10]: if section != current_section:
emoji = p.get_status_emoji() text += f"{section}:\n"
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" 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} ({suffix})",
callback_data=f"tdetail_{tid}"
))
text += ( self._add_nav(builder, "streamers_page_", page, total_pages)
f"{emoji} {p.channel}\n"
f" 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n"
)
builder.row(InlineKeyboardButton(
text=f"{emoji} {p.channel} (🔗{p.links_found})",
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}"
))
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
await self._edit_or_send(message, text, builder.as_markup(), edit) 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() active = await self.task_manager.get_active_tasks()
completed = await self.task_manager.get_completed_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")} def _pick(src, types):
visit_done = {tid: p for tid, p in completed.items() if p.task_type in ("visit", "user_visit")} return [(tid, p) for tid, p in src.items() if p.task_type in types]
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 = [("📺 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() builder = InlineKeyboardBuilder()
if not any([visit_active, visit_done, twitch_active, twitch_done]): if not entries:
text = "📊 Задачи\n\nЗадач нет." text = "📊 Задачи\n\nЗадач нет."
else: 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: current_section = None
text += "📺 Twitch — активные:\n" for section, tid, p in page_items:
for tid, p in list(twitch_active.items())[:10]: if section != current_section:
em = p.get_status_emoji() text += f"\n{section}:\n" if current_section else f"{section}:\n"
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" current_section = section
text += f"{em} {p.channel}\n 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n" em = p.get_status_emoji()
builder.row(InlineKeyboardButton( is_active = "активные" in section
text=f"{em} 📺 {p.channel} (🔗{p.links_found})", if section.startswith("📺"):
callback_data=f"tdetail_{tid}" 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"
if twitch_done: label = f"{em} 📺 {p.channel} (🔗{p.links_found})"
text += "\n📺 Twitch — завершённые:\n" else:
for tid, p in list(twitch_done.items())[:5]: text += f"{em} {p.channel}\n"
em = p.get_status_emoji() label = f"{em} 📺 {p.channel} (завершена)"
text += f"{em} {p.channel}\n" else:
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()
url_s = (p.url or "")[:35] url_s = (p.url or "")[:35]
uid_tag = f" | uid:{p.user_id}" if p.user_id else "" uid_tag = f" | uid:{p.user_id}" if p.user_id else ""
text += f"{em} {url_s}{uid_tag}\n{p.successful_visits}/{p.max_visits} | 📊{p.total_visits}\n" if is_active:
builder.row(InlineKeyboardButton( text += f"{em} {url_s}{uid_tag}\n{p.successful_visits}/{p.max_visits} | 📊{p.total_visits}\n"
text=f"{em} 🔗 {url_s[:28]} ({p.successful_visits}/{p.max_visits})", label = f"{em} 🔗 {url_s[:28]} ({p.successful_visits}/{p.max_visits})"
callback_data=f"tdetail_{tid}" else:
)) text += f"{em} {url_s}{uid_tag}\n"
label = f"{em} 🔗 {url_s[:28]} (завершена)"
builder.row(InlineKeyboardButton(text=label, callback_data=f"tdetail_{tid}"))
if visit_done: self._add_nav(builder, "tasks_page_", page, total_pages)
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 ""
text += f"{em} {url_s}{uid_tag}\n"
builder.row(InlineKeyboardButton(
text=f"{em} 🔗 {url_s[:28]} (завершена)",
callback_data=f"tdetail_{tid}"
))
builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_tasks")) builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_tasks"))
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main")) builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
@@ -2125,8 +2147,29 @@ class BotInterface:
reply_markup=builder.as_markup() 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, 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 was_offline = params.stream_offline
@@ -2161,6 +2204,10 @@ class BotInterface:
elif viewers > 0 and was_offline: elif viewers > 0 and was_offline:
was_offline = False was_offline = False
# Началась новая трансляция — чистим очередь прошлой.
# Делаем это до снятия паузы, иначе отложенные серии успеют
# проснуться и отработать по неактуальным ссылкам.
dropped = await self._drop_pending_series(params, url_tasks)
params.stream_offline = False params.stream_offline = False
if params.auto_paused: if params.auto_paused:
params.paused = False params.paused = False
@@ -2168,11 +2215,14 @@ class BotInterface:
await self.task_manager.update_task(task_id, paused=False) await self.task_manager.update_task(task_id, paused=False)
if self.storage: if self.storage:
await self.storage.save_task(task_id, params) await self.storage.save_task(task_id, params)
logger.info(f"[watcher] {params.channel} is online ({viewers} viewers)") logger.info(
msg = await send_message_safe( f"[watcher] {params.channel} is online ({viewers} viewers), "
message.bot, params.chat_id, f"queue cleared: {dropped} clicks dropped"
f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён"
) )
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: if msg and params.user_id:
await self._track_msg(params.user_id, msg.message_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}") 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: try:
await irc.listen_for_messages( 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): async def _process_url(self, url: str, username: str, task_id: str, params: TaskParams, message: Message):
"""Обрабатывает найденную ссылку (выполняется параллельно).""" """Обрабатывает найденную ссылку (выполняется параллельно)."""
notif_msg = None
try: try:
params.links_found += 1 params.links_found += 1
@@ -2389,7 +2442,7 @@ class BotInterface:
params.pending_visits += visits_count params.pending_visits += visits_count
params.total_planned_visits += visits_count params.total_planned_visits += visits_count
logger.info(f"🔗 @{username} {url}{visits_count} clicks ({calc_info}), total pending: {params.pending_visits}") 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, message.bot, message.chat.id,
f"🔗 @{username}: {url[:60]}\n" f"🔗 @{username}: {url[:60]}\n"
f"📊 {calc_info}\n" f"📊 {calc_info}\n"
@@ -2402,6 +2455,7 @@ class BotInterface:
# Resume > 5 мин: серия этой ссылки полностью отбрасывается # Resume > 5 мин: серия этой ссылки полностью отбрасывается
SERIES_EXPIRY = 5 * 60 SERIES_EXPIRY = 5 * 60
remaining = visits_count remaining = visits_count
successful_this_link = 0
paused_at: Optional[float] = None paused_at: Optional[float] = None
first_click = True first_click = True
@@ -2443,6 +2497,7 @@ class BotInterface:
return return
if result.success: if result.success:
params.successful_visits += 1 params.successful_visits += 1
successful_this_link += 1
remaining -= 1 remaining -= 1
if params.user_id: if params.user_id:
new_balance = await self.balance_storage.deduct(params.user_id, 1) new_balance = await self.balance_storage.deduct(params.user_id, 1)
@@ -2470,9 +2525,23 @@ class BotInterface:
if params.stopped: if params.stopped:
return return
except asyncio.CancelledError:
# Очередь сброшена (началась новая трансляция) — серия прекращается
logger.info(f"🧹 Серия отменена: {url[:40]}")
raise
except Exception as e: except Exception as e:
logger.error(f"Process URL error: {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
# ========================================================================= # =========================================================================
# КАБИНЕТ ПОЛЬЗОВАТЕЛЯ # КАБИНЕТ ПОЛЬЗОВАТЕЛЯ
@@ -2646,19 +2715,7 @@ class BotInterface:
if not re.match(r'^[a-zA-Z0-9_]+$', text): if not re.match(r'^[a-zA-Z0-9_]+$', text):
await message.answer("❌ Неверное название канала (только буквы, цифры, _). Попробуйте ещё раз:") await message.answer("❌ Неверное название канала (только буквы, цифры, _). Попробуйте ещё раз:")
return return
channel = text.lower() 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 # Удаляем промпт шага 1 перед отправкой шага 2
if user_id in self._menu_msg: if user_id in self._menu_msg:
await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id)) await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id))
@@ -3069,51 +3126,76 @@ class BotInterface:
else: else:
await self._send_temp(message, "Ошибка запуска задачи") await self._send_temp(message, "Ошибка запуска задачи")
async def _show_users_list(self, message: Message, edit: bool = False): async def _show_users_list(self, message: Message, edit: bool = False, page: int = 0):
balances = await self.balance_storage.get_all() profiles, all_tasks, all_clicks, all_rub = await asyncio.gather(
all_tasks = await self.task_manager.get_all_tasks() 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() builder = InlineKeyboardBuilder()
if not balances: if not all_user_ids:
text += "Нет пользователей с балансом." text = "👥 Пользователи (0)\n\nНет зарегистрированных пользователей."
else: else:
for uid_str, balance in list(balances.items())[:15]: page_ids, page, total_pages = self._paginate(all_user_ids, page)
uid = int(uid_str) 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( active = sum(
1 for p in all_tasks.values() 1 for p in all_tasks.values()
if p.user_id == uid and not p.completed and not p.stopped 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( builder.row(InlineKeyboardButton(
text=f"👤 {uid} (💰{balance})", text=f"👤 {name or uid} {bal_str}",
callback_data=f"udetail_{uid}", callback_data=f"udetail_{uid}",
)) ))
self._add_nav(builder, "users_page_", page, total_pages)
builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main")) builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main"))
await self._edit_or_send(message, text, builder.as_markup(), edit) await self._edit_or_send(message, text, builder.as_markup(), edit)
async def _show_user_detail(self, callback: CallbackQuery, target_user_id: int): 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() 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} 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) 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 = ( text = (
f"👤 Пользователь {target_user_id}\n\n" f"{header}\n\n"
f"🔢 Переходов: {balance}\n" f"🖱 Переходов: {clicks}\n"
f"📋 Задач: {len(user_tasks)} (активных: {active})\n\n" f"🪙 Рублей: {rub}\n"
"Пополнить баланс:" f"📋 Задач: {len(user_tasks)} (активных: {active})\n"
) )
builder = InlineKeyboardBuilder() builder = InlineKeyboardBuilder()
for amount in [10, 50, 100, 500]:
builder.button(text=f"+{amount}", callback_data=f"ubal_{amount}_{target_user_id}") builder.row(InlineKeyboardButton(text=" Переходы", callback_data=f"ubal_menu_{target_user_id}"))
builder.adjust(4) builder.row(InlineKeyboardButton(text=" Рубли", callback_data=f"urub_menu_{target_user_id}"))
builder.row(InlineKeyboardButton(
text="✏️ Другая сумма",
callback_data=f"ubal_custom_{target_user_id}",
))
if user_tasks: if user_tasks:
builder.row(InlineKeyboardButton( builder.row(InlineKeyboardButton(
text=f"📋 Задачи ({len(user_tasks)})", text=f"📋 Задачи ({len(user_tasks)})",
+42 -2
View File
@@ -6,7 +6,7 @@ import asyncio
import json import json
import logging import logging
from pathlib import Path from pathlib import Path
from datetime import datetime from datetime import datetime, timedelta
from typing import Dict, Optional from typing import Dict, Optional
from managers.task_manager import TaskParams from managers.task_manager import TaskParams
@@ -309,6 +309,9 @@ class RubleBalanceStorage(_IntBalanceStorage):
class PaymentStorage: class PaymentStorage:
"""Хранилище ожидающих платежей.""" """Хранилище ожидающих платежей."""
# Сколько держать обработанные платежи ради защиты от дублей вебхука
PROCESSED_TTL_DAYS = 7
def __init__(self, file_path: str = "data/payments.json"): def __init__(self, file_path: str = "data/payments.json"):
self.file_path = Path(file_path) self.file_path = Path(file_path)
self.file_path.parent.mkdir(parents=True, exist_ok=True) 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: with open(self.file_path, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2) 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 def save(self, payment_id: str, data: dict) -> None:
async with self._lock: async with self._lock:
payments = await asyncio.to_thread(self._load_sync) payments = await asyncio.to_thread(self._load_sync)
payments = self._prune_processed(payments)
payments[payment_id] = data payments[payment_id] = data
await asyncio.to_thread(self._save_sync, payments) await asyncio.to_thread(self._save_sync, payments)
@@ -344,11 +367,28 @@ class PaymentStorage:
payments.pop(payment_id, None) payments.pop(payment_id, None)
await asyncio.to_thread(self._save_sync, payments) 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 def get_by_user(self, user_id: int) -> Optional[dict]:
"""Возвращает незакрытый счёт пользователя.
Обработанные платежи остаются в файле как защита от повторного
начисления по дублю вебхука, но открытым счётом уже не считаются.
"""
async with self._lock: async with self._lock:
payments = await asyncio.to_thread(self._load_sync) payments = await asyncio.to_thread(self._load_sync)
for pid, data in payments.items(): 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 {**data, "payment_id": pid}
return None return None
+1 -1
View File
@@ -332,7 +332,7 @@ class BrowserService:
current = new_url current = new_url
redirect = True redirect = True
task = asyncio.create_task( 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 self._move_mouse(page, task)
await task await task
+12 -1
View File
@@ -5,6 +5,7 @@ Twitch чат через IRC over WebSocket (wss://irc-ws.chat.twitch.tv:443).
import asyncio import asyncio
import logging import logging
import random
import re import re
from typing import Optional, Callable, Awaitable, List from typing import Optional, Callable, Awaitable, List
from datetime import datetime from datetime import datetime
@@ -34,7 +35,13 @@ class TwitchGQLChatClient:
): ):
self.channel = channel.lower() self.channel = channel.lower()
self.target_username = target_username.lower() self.target_username = target_username.lower()
self.irc_username = irc_username or settings.IRC_USERNAME # 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.irc_oauth = irc_oauth or settings.IRC_OAUTH
self._ws: Optional[aiohttp.ClientWebSocketResponse] = None self._ws: Optional[aiohttp.ClientWebSocketResponse] = None
@@ -257,6 +264,10 @@ class TwitchGQLChatClient:
raise raise
except Exception: except Exception:
pass pass
else:
logger.info(
f"🚫 Домен не разрешён: {url} (allowed_domains={allowed_domains})"
)
return stats return stats
+8 -1
View File
@@ -3,6 +3,7 @@ IRC сервис для мониторинга Twitch чата.
""" """
import asyncio import asyncio
import random
import re import re
import logging import logging
from typing import Optional, Callable, Awaitable, List from typing import Optional, Callable, Awaitable, List
@@ -28,7 +29,13 @@ class TwitchIRCClient:
): ):
self.channel = channel.lower() self.channel = channel.lower()
self.target_username = target_username.lower() self.target_username = target_username.lower()
self.irc_username = irc_username or settings.IRC_USERNAME # 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.irc_oauth = irc_oauth or settings.IRC_OAUTH
self._reader: Optional[asyncio.StreamReader] = None self._reader: Optional[asyncio.StreamReader] = None
+24 -3
View File
@@ -11,6 +11,7 @@ import hashlib
import json import json
import logging import logging
from datetime import datetime from datetime import datetime
from pathlib import Path
from typing import Optional from typing import Optional
from aiohttp import web from aiohttp import web
@@ -74,6 +75,21 @@ class HelketWebhookServer:
await self._runner.cleanup() await self._runner.cleanup()
logger.info("Heleket webhook server stopped") 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: async def _handle_health(self, request: web.Request) -> web.Response:
return web.Response(text="ok") return web.Response(text="ok")
@@ -102,9 +118,15 @@ class HelketWebhookServer:
logger.info(f"Webhook: uuid={payment_uuid} status={status} — ignored") logger.info(f"Webhook: uuid={payment_uuid} status={status} — ignored")
return web.Response(text="ok") 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: 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") return web.Response(text="ok")
user_id = payment.get("user_id") user_id = payment.get("user_id")
@@ -147,7 +169,6 @@ class HelketWebhookServer:
} }
logger.info(f"Webhook: user={user_id} +{rub} RUB") 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) await self._history_storage.add(user_id, history_record)
if chat_id: if chat_id: