Compare commits

..

16 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
Yuriy Yuriev f262616aad fix click 2026-06-05 20:56:26 +07:00
Yuriy Yuriev e438f09693 fix delete old state 2026-05-29 22:06:34 +07:00
Yuriy Yuriev b02b7518cf add stiky 2026-05-29 21:47:40 +07:00
Yuriy Yuriev 61ff90e51b fix to irc 2026-05-26 18:04:52 +07:00
Yuriy Yuriev 493c67cb1d fix gql 2026-05-26 18:01:07 +07:00
Yuriy Yuriev 795b2e928c fix url 2026-05-26 17:58:47 +07:00
Yuriy Yuriev d95dc462d9 Add graphQl for chat 2026-05-26 17:52:40 +07:00
Yuriy Yuriev 406a58c585 add restart storage 2026-05-23 00:07:44 +07:00
Yuriy Yuriev 68275386bf add message user 2026-05-22 23:24:29 +07:00
Yuriy Yuriev e2e2761b51 fix comand 2026-05-22 23:06:31 +07:00
15 changed files with 1444 additions and 546 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'\\)\")"
] ]
} }
} }
+46 -12
View File
@@ -1,8 +1,8 @@
import asyncio import json
import hashlib
import secrets
import logging import logging
import secrets
from datetime import datetime, timedelta from datetime import datetime, timedelta
from pathlib import Path
from typing import Dict, Optional from typing import Dict, Optional
from config.settings import settings from config.settings import settings
@@ -10,10 +10,15 @@ from config.settings import settings
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
ADMIN_PASSWORD = settings.ADMIN_PASSWORD ADMIN_PASSWORD = settings.ADMIN_PASSWORD
_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] = {}
@@ -22,21 +27,27 @@ class AuthManager:
self._max_attempts = 5 self._max_attempts = 5
self._lockout_duration = 300 self._lockout_duration = 300
self._session_timeout = timedelta(minutes=session_timeout_minutes) self._session_timeout = timedelta(minutes=session_timeout_minutes)
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:
return False return False
# Обычные пользователи — сессия бесконечная
if self._roles.get(user_id) != "admin":
return True
# Админы — таймаут применяется
if datetime.now() - self._sessions[user_id] > self._session_timeout: if datetime.now() - self._sessions[user_id] > self._session_timeout:
self.logout(user_id) self.logout(user_id)
return False return False
@@ -49,14 +60,37 @@ class AuthManager:
self._sessions[user_id] = datetime.now() self._sessions[user_id] = datetime.now()
self._roles[user_id] = "admin" if is_admin else "user" self._roles[user_id] = "admin" if is_admin else "user"
self._failed_attempts.pop(user_id, None) self._failed_attempts.pop(user_id, None)
self._save_sessions()
logger.info(f"User {user_id} logged in as {self._roles[user_id]}") logger.info(f"User {user_id} logged in as {self._roles[user_id]}")
def logout(self, user_id: int): def logout(self, user_id: int):
self._sessions.pop(user_id, None) self._sessions.pop(user_id, None)
self._roles.pop(user_id, None) self._roles.pop(user_id, None)
self._failed_attempts.pop(user_id, None) self._failed_attempts.pop(user_id, None)
self._save_sessions()
logger.info(f"User {user_id} logged out") logger.info(f"User {user_id} logged out")
def _save_sessions(self) -> None:
try:
_SESSIONS_FILE.parent.mkdir(parents=True, exist_ok=True)
data = {str(uid): role for uid, role in self._roles.items()}
_SESSIONS_FILE.write_text(json.dumps(data), encoding="utf-8")
except Exception as e:
logger.warning(f"Failed to save sessions: {e}")
def _restore_sessions(self) -> None:
try:
if not _SESSIONS_FILE.exists():
return
data = json.loads(_SESSIONS_FILE.read_text(encoding="utf-8"))
for uid_str, role in data.items():
uid = int(uid_str)
self._sessions[uid] = datetime.now()
self._roles[uid] = role
logger.info(f"Restored {len(data)} sessions")
except Exception as e:
logger.warning(f"Failed to restore sessions: {e}")
def is_locked_out(self, user_id: int) -> bool: def is_locked_out(self, user_id: int) -> bool:
if user_id not in self._failed_attempts: if user_id not in self._failed_attempts:
return False return False
+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 {}
+6 -5
View File
@@ -54,11 +54,12 @@ class Settings(BaseSettings):
SOCKET_TIMEOUT: int = 30 SOCKET_TIMEOUT: int = 30
RECV_TIMEOUT: float = 5.0 RECV_TIMEOUT: float = 5.0
# Proxy # Sticky proxy (PlainProxies / аналоги)
PROXY_FILE: str = "proxies.txt" # При заполненном STICKY_PROXY_HOST rotating-прокси для браузера не используются
PROXY_COOLDOWN_TIME: int = 30 STICKY_PROXY_HOST: str = " " # res-unlimited-XXXX.plainproxies.com
PROXY_MAX_USAGE_BEFORE_COOLDOWN: int = 5 STICKY_PROXY_PORT: int = 8080
STICKY_PROXY_USER: str = " " # логин из личного кабинета
STICKY_PROXY_PASS: str = " " # пароль из личного кабинета
# Screenshots # Screenshots
SCREENSHOTS_DIR: str = "" SCREENSHOTS_DIR: str = ""
+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)
+445 -345
View File
File diff suppressed because it is too large Load Diff
+41 -30
View File
@@ -10,7 +10,6 @@ from aiogram.types import BotCommand
from aiogram.client.session.aiohttp import AiohttpSession from aiogram.client.session.aiohttp import AiohttpSession
from config.settings import settings from config.settings import settings
from core.logger import setup_logger from core.logger import setup_logger
from managers.proxy_manager import ProxyManager
from managers.background_tasks import BackgroundTaskManager from managers.background_tasks import BackgroundTaskManager
from services.browser_service import BrowserService from services.browser_service import BrowserService
from services.browser_pool import BrowserPool from services.browser_pool import BrowserPool
@@ -43,10 +42,10 @@ class BotApplication:
def __init__(self): def __init__(self):
self.bot: Bot = None self.bot: Bot = None
self.dispatcher: Dispatcher = None self.dispatcher: Dispatcher = None
self._initialized = False # True только после успешного _restore_tasks
self.proxy_manager = ProxyManager()
self.background_tasks = BackgroundTaskManager() self.background_tasks = BackgroundTaskManager()
self.browser_service = BrowserService(self.proxy_manager) self.browser_service = BrowserService()
self.browser_pool = BrowserPool(settings.BROWSER_POOL_SIZE, self.browser_service) self.browser_pool = BrowserPool(settings.BROWSER_POOL_SIZE, self.browser_service)
self.storage = TaskStorage() self.storage = TaskStorage()
self.chat_storage = ChatStorage() self.chat_storage = ChatStorage()
@@ -55,7 +54,6 @@ class BotApplication:
self.interface = BotInterface( self.interface = BotInterface(
background_tasks=self.background_tasks, background_tasks=self.background_tasks,
proxy_manager=self.proxy_manager,
browser_service=self.browser_service, browser_service=self.browser_service,
browser_pool=self.browser_pool, browser_pool=self.browser_pool,
storage=self.storage, storage=self.storage,
@@ -85,9 +83,6 @@ class BotApplication:
self.interface.register(self.dispatcher) self.interface.register(self.dispatcher)
await self._update_bot_commands() await self._update_bot_commands()
# Определяем типы прокси без явного протокола
await self.proxy_manager.detect_types()
# Запускаем пул браузеров # Запускаем пул браузеров
await self.browser_pool.start() await self.browser_pool.start()
logger.info(f"Browser pool: {settings.BROWSER_POOL_SIZE} workers") logger.info(f"Browser pool: {settings.BROWSER_POOL_SIZE} workers")
@@ -107,31 +102,46 @@ class BotApplication:
# Восстанавливаем задачи # Восстанавливаем задачи
await self._restore_tasks() await self._restore_tasks()
logger.info(f"Proxies: {self.proxy_manager.count}") self._initialized = True
# Периодическое сохранение на случай аварийного завершения
asyncio.create_task(self._autosave_loop(), name="autosave")
logger.info("Bot initialized!") logger.info("Bot initialized!")
async def _autosave_loop(self):
"""Сохраняет состояние задач каждые 60 секунд."""
while True:
await asyncio.sleep(60)
try:
tasks = dict(self.interface.task_manager._tasks)
await self.storage.save_tasks(tasks)
logger.debug(f"Autosaved {len(tasks)} tasks")
except Exception as e:
logger.warning(f"Autosave failed: {e}")
async def _restore_tasks(self): async def _restore_tasks(self):
"""Восстанавливает задачи после перезапуска.""" """Восстанавливает задачи после перезапуска."""
saved_tasks = await self.storage.load_tasks() saved_tasks = await self.storage.load_tasks()
if not saved_tasks: if not saved_tasks:
logger.info("No tasks to restore") logger.info("No tasks to restore")
return return
restored = 0 restored = 0
notify_users: dict = {} # user_id -> chat_id # user_id -> (chat_id, was_paused)
notify_users: dict = {}
for task_id, params in saved_tasks.items(): for task_id, params in saved_tasks.items():
# Завершённые и остановленные — загружаем в память но не запускаем # Завершённые и остановленные — только в память, не запускаем
if params.completed or params.stopped: if params.completed or params.stopped:
await self.interface.task_manager.add_task(task_id, params) await self.interface.task_manager.add_task(task_id, params)
continue continue
# Активные задачи ставим на паузу и перезапускаем was_paused = params.paused # сохраняем намерение пользователя
# Если стрим был оффлайн — сохраняем этот статус
params.paused = True
if params.stream_offline: if params.stream_offline:
logger.info(f"Task {task_id}: stream was offline, keeping paused") logger.info(f"Task {task_id}: stream was offline, keeping paused")
await self.interface.task_manager.add_task(task_id, params) await self.interface.task_manager.add_task(task_id, params)
if params.task_type == "twitch_irc" and params.chat_id: if params.task_type == "twitch_irc" and params.chat_id:
@@ -156,20 +166,19 @@ class BotApplication:
restored += 1 restored += 1
if params.user_id and params.chat_id: if params.user_id and params.chat_id:
notify_users[params.user_id] = params.chat_id notify_users[params.user_id] = (params.chat_id, was_paused)
# Уведомляем пользователей о рестарте # Уведомляем пользователей разные сообщения для паузы и активных задач
for uid, chat_id in notify_users.items(): for uid, (chat_id, was_paused) in notify_users.items():
try: try:
await self.bot.send_message( if was_paused:
chat_id, msg = "🔄 Бот был перезапущен\n\n⏸ Ваша задача остаётся на паузе."
"🔄 Бот был перезапущен\n\n" else:
"Ваши задачи приостановлены.\n" msg = "🔄 Бот был перезапущен\n\n▶️ Ваши задачи продолжают работу автоматически."
"Нажмите ▶️ в списке задач для возобновления." await self.bot.send_message(chat_id, msg)
)
except Exception as e: except Exception as e:
logger.warning(f"Failed to notify user {uid}: {e}") logger.warning(f"Failed to notify user {uid}: {e}")
logger.info(f"🔄 Restored {restored} tasks") logger.info(f"🔄 Restored {restored} tasks")
async def start(self): async def start(self):
@@ -181,10 +190,12 @@ class BotApplication:
async def stop(self): async def stop(self):
"""Остановка с сохранением.""" """Остановка с сохранением."""
logger.info("Shutting down...") logger.info("Shutting down...")
# Сохраняем задачи # Сохраняем только если инициализация прошла успешно —
tasks = self.interface.task_manager._tasks # иначе task_manager пустой и мы затрём уже сохранённые задачи
await self.storage.save_tasks(tasks) if self._initialized:
tasks = self.interface.task_manager._tasks
await self.storage.save_tasks(tasks)
await self.background_tasks.cancel_all() await self.background_tasks.cancel_all()
await self.browser_pool.stop() await self.browser_pool.stop()
+59 -9
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
@@ -54,6 +54,7 @@ class TaskStorage:
"max_ctr": params.max_ctr, "max_ctr": params.max_ctr,
"started_at": params.started_at.isoformat() if params.started_at else None, "started_at": params.started_at.isoformat() if params.started_at else None,
"stream_offline": params.stream_offline, "stream_offline": params.stream_offline,
"remaining_clicks": params.remaining_clicks,
# pending_visits и total_planned_visits — runtime, не сохраняем # pending_visits и total_planned_visits — runtime, не сохраняем
} }
@@ -88,6 +89,7 @@ class TaskStorage:
min_ctr=data.get("min_ctr", 0.8), min_ctr=data.get("min_ctr", 0.8),
max_ctr=data.get("max_ctr", 1.0), max_ctr=data.get("max_ctr", 1.0),
stream_offline=data.get("stream_offline", False), stream_offline=data.get("stream_offline", False),
remaining_clicks=data.get("remaining_clicks", 0),
) )
if data.get("started_at"): if data.get("started_at"):
try: try:
@@ -133,15 +135,23 @@ class TaskStorage:
return {} return {}
async def save_task(self, task_id: str, params: TaskParams) -> None: async def save_task(self, task_id: str, params: TaskParams) -> None:
tasks = await self.load_tasks() async with self._lock:
tasks[task_id] = params try:
await self.save_tasks(tasks) data = await asyncio.to_thread(self._load_sync)
data[task_id] = self._params_to_dict(params)
await asyncio.to_thread(self._save_sync, data)
except Exception as e:
logger.error(f"save_task error: {e}")
async def delete_task(self, task_id: str) -> None: async def delete_task(self, task_id: str) -> None:
tasks = await self.load_tasks() async with self._lock:
if task_id in tasks: try:
del tasks[task_id] data = await asyncio.to_thread(self._load_sync)
await self.save_tasks(tasks) if task_id in data:
del data[task_id]
await asyncio.to_thread(self._save_sync, data)
except Exception as e:
logger.error(f"delete_task error: {e}")
class ChatStorage: class ChatStorage:
@@ -299,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)
@@ -317,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)
@@ -334,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
+19 -8
View File
@@ -59,6 +59,7 @@ class TaskParams:
links_found: int = 0 links_found: int = 0
pending_visits: int = 0 # кликов в очереди (runtime, уменьшается) pending_visits: int = 0 # кликов в очереди (runtime, уменьшается)
total_planned_visits: int = 0 # всего запланировано кликов по всем ссылкам (runtime) total_planned_visits: int = 0 # всего запланировано кликов по всем ссылкам (runtime)
remaining_clicks: int = 0 # остаток кликов в текущей серии (сохраняется)
chat_id: Optional[int] = None chat_id: Optional[int] = None
user_id: Optional[int] = None # Telegram user_id, если задача создана пользователем user_id: Optional[int] = None # Telegram user_id, если задача создана пользователем
@@ -271,8 +272,12 @@ class TaskManager:
return await self.get_tasks_by_type("twitch_irc") return await self.get_tasks_by_type("twitch_irc")
async def get_visit_tasks(self) -> Dict[str, TaskParams]: async def get_visit_tasks(self) -> Dict[str, TaskParams]:
"""Только задачи посещения.""" """Задачи посещения (admin и user)."""
return await self.get_tasks_by_type("visit") async with self._lock:
return {
tid: p for tid, p in self._tasks.items()
if p.task_type in ("visit", "user_visit")
}
async def get_stats(self) -> dict: async def get_stats(self) -> dict:
"""Общая статистика.""" """Общая статистика."""
@@ -280,7 +285,7 @@ class TaskManager:
for p in self._tasks.values(): for p in self._tasks.values():
if p.task_type == "twitch_irc": if p.task_type == "twitch_irc":
twitch += 1 twitch += 1
elif p.task_type == "visit": elif p.task_type in ("visit", "user_visit"):
visit += 1 visit += 1
if p.paused: if p.paused:
paused += 1 paused += 1
@@ -297,14 +302,20 @@ class TaskManager:
def _get_task_info(self, task_id: str, params: TaskParams) -> str: def _get_task_info(self, task_id: str, params: TaskParams) -> str:
"""Информация о задаче для отображения.""" """Информация о задаче для отображения."""
status = params.get_status_emoji() status = params.get_status_emoji()
if params.task_type == "twitch_irc": if params.task_type == "twitch_irc":
user_tag = f" | 👤uid:{params.user_id}" if params.user_id else ""
return ( return (
f"{status} 📺 **{params.channel}**\n" f"{status} 📺 **{params.channel}**\n"
f" `{task_id[:12]}...`\n" f" `{task_id[:12]}...`{user_tag}\n"
f" 👤 @{params.target_username} | " f" 🔗{params.links_found} | 📊{params.total_visits}\n"
f"🔗{params.links_found} | " )
f"📊{params.total_visits}\n" elif params.task_type == "user_visit":
user_tag = f" | 👤uid:{params.user_id}" if params.user_id else ""
return (
f"{status} 🔗 `{params.url[:40]}`\n"
f" `{task_id[:12]}...`{user_tag}\n"
f"{params.successful_visits}/{params.max_visits} | 📊{params.total_visits}\n"
) )
else: else:
return ( return (
+86 -86
View File
@@ -1,5 +1,5 @@
""" """
Browser service с поддержкой SOCKS5 через локальный HTTP туннель. Browser service с поддержкой HTTP и sticky-прокси.
""" """
import asyncio import asyncio
@@ -12,11 +12,9 @@ from dataclasses import dataclass
from camoufox import DefaultAddons from camoufox import DefaultAddons
from camoufox.async_api import AsyncCamoufox from camoufox.async_api import AsyncCamoufox
from camoufox.exceptions import InvalidIP from camoufox.exceptions import InvalidIP, InvalidProxy
from config.settings import settings from config.settings import settings
from managers.proxy_manager import ProxyManager, Proxy, ProxyType
from services.socks5_to_http_proxy import Socks5ProxyPool
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -73,6 +71,44 @@ def _patch_camoufox():
_patch_camoufox() _patch_camoufox()
def _make_sticky_proxy_config() -> tuple:
"""Возвращает (proxy_config, proxy_url_with_auth, proxy_info) для sticky-сессии."""
import re
session_id = random.randint(100000, 999999)
# Заменяем session-XXXXXX в username, если уже есть — иначе добавляем
base = settings.STICKY_PROXY_USER
if re.search(r'session-\d+', base):
username = re.sub(r'session-\d+', f'session-{session_id}', base)
else:
username = f"{base}-session-{session_id}"
config = {
"server": f"http://{settings.STICKY_PROXY_HOST}:{settings.STICKY_PROXY_PORT}",
"username": username,
"password": settings.STICKY_PROXY_PASS,
}
proxy_url = (
f"http://{username}:{settings.STICKY_PROXY_PASS}"
f"@{settings.STICKY_PROXY_HOST}:{settings.STICKY_PROXY_PORT}"
)
return config, proxy_url, f"sticky-{session_id}@{settings.STICKY_PROXY_HOST}"
async def _resolve_sticky_ip(proxy_url: str) -> Optional[str]:
"""Определяет реальный IP sticky-сессии через сам прокси."""
import aiohttp
try:
async with aiohttp.ClientSession() as session:
async with session.get(
"https://api.ipify.org",
proxy=proxy_url,
timeout=aiohttp.ClientTimeout(total=10),
) as resp:
return (await resp.text()).strip()
except Exception as e:
logger.warning(f"Sticky proxy IP resolve failed: {e}")
return None
def _locale_for_ip(ip: str) -> str: def _locale_for_ip(ip: str) -> str:
"""Возвращает BCP47 locale без script-тега по реальному IP прокси.""" """Возвращает BCP47 locale без script-тега по реальному IP прокси."""
try: try:
@@ -104,14 +140,15 @@ class VisitResult:
class BrowserService: class BrowserService:
""" """
Сервис для посещения страниц через Camoufox. Сервис для посещения страниц через Camoufox.
Поддерживает HTTP, SOCKS5 (через туннель) и прямое соединение. Поддерживает HTTP и sticky-прокси.
""" """
def __init__(self, proxy_manager: ProxyManager): _PROXY_FAIL_THRESHOLD = 5 # пауза задачи после N подряд неудачных визитов
self.proxy_manager = proxy_manager
self.socks5_pool = Socks5ProxyPool(idle_timeout=300) def __init__(self):
self.screenshots_dir = Path(settings.SCREENSHOTS_DIR) self.screenshots_dir = Path(settings.SCREENSHOTS_DIR)
self.screenshots_dir.mkdir(exist_ok=True) self.screenshots_dir.mkdir(exist_ok=True)
self._proxy_fail_streak: int = 0
async def visit_page( async def visit_page(
self, self,
@@ -125,101 +162,70 @@ class BrowserService:
settings.DEFAULT_MAX_READING settings.DEFAULT_MAX_READING
) )
if not settings.STICKY_PROXY_HOST:
return VisitResult(url=url, success=False, error="No proxy configured")
# Максимальное время одной попытки: запуск браузера + загрузка страницы + чтение + буфер
attempt_timeout = reading_time + 150
for attempt in range(3): for attempt in range(3):
proxy = None
use_socks5 = False
try: try:
if self.proxy_manager.has_proxies(): proxy_config, proxy_url, proxy_info = _make_sticky_proxy_config()
proxy = await self.proxy_manager.acquire_proxy(timeout=30) real_ip = await _resolve_sticky_ip(proxy_url)
if real_ip is None:
logger.warning(f"Sticky proxy IP resolve failed (attempt {attempt + 1}), retrying with new session")
continue
_set_real_ip(real_ip)
proxy_locale = _locale_for_ip(real_ip)
logger.info(f"Using sticky proxy: {proxy_info} ip={real_ip}")
camoufox_kwargs = dict(
headless=True,
geoip=True,
humanize=True,
exclude_addons=[DefaultAddons.UBO],
proxy=proxy_config,
)
if proxy_locale:
camoufox_kwargs["locale"] = proxy_locale
proxy_locale = None async def _run_browser():
try:
if proxy and proxy.proxy_type in [ProxyType.SOCKS5, ProxyType.SOCKS4]:
proxy_config = await self.socks5_pool.get_proxy_config(proxy)
use_socks5 = True
real_ip = proxy.id.split("://")[-1].split(":")[0]
_set_real_ip(real_ip)
proxy_locale = _locale_for_ip(real_ip)
has_auth = bool(proxy.login and proxy.password)
logger.info(f"Using SOCKS5 tunnel for: {proxy.id} auth={'yes' if has_auth else 'NO'} locale={proxy_locale}")
elif proxy:
proxy_config = proxy.proxy_config
real_ip = proxy.id.split("://")[-1].split(":")[0]
_set_real_ip(real_ip)
proxy_locale = _locale_for_ip(real_ip)
logger.info(f"Using HTTP proxy: {proxy.id} locale={proxy_locale}")
else:
proxy_config = None
_set_real_ip(None)
logger.info("Direct connection")
except asyncio.CancelledError:
# Задача отменена — прокси не виноват, освобождаем без штрафа
if proxy:
await self.proxy_manager.release_proxy(proxy, success=True)
raise
if proxy_config:
camoufox_kwargs = dict(
headless=True,
geoip=True,
humanize=True,
exclude_addons=[DefaultAddons.UBO],
proxy=proxy_config,
)
if proxy_locale:
camoufox_kwargs["locale"] = proxy_locale
async with AsyncCamoufox(**camoufox_kwargs) as browser: async with AsyncCamoufox(**camoufox_kwargs) as browser:
result = await self._browse_page( return await self._browse_page(browser, url, proxy_info, reading_time)
browser, url,
proxy.id if proxy else "direct", result = await asyncio.wait_for(_run_browser(), timeout=attempt_timeout)
reading_time
)
else:
return VisitResult(url=url, success=False, error="No proxy available")
# Определяем тип ошибки: ошибка соединения (прокси сломан) vs ошибка загрузки (прокси ok)
is_proxy_dead = result.error and any( is_proxy_dead = result.error and any(
k in result.error for k in ("502", "Bad Gateway", "SOCKS", "NS_ERROR_PROXY") k in result.error for k in ("502", "Bad Gateway", "NS_ERROR_PROXY")
) )
is_load_error = result.error and not is_proxy_dead and any( is_load_error = result.error and not is_proxy_dead and any(
k in result.error for k in ("Timeout", "ERR_", "NS_ERROR_", "Connection") k in result.error for k in ("Timeout", "ERR_", "NS_ERROR_", "Connection")
) )
if proxy:
# Таймаут/ошибка загрузки — прокси не виноват, освобождаем без штрафа
release_success = True if is_load_error else result.success
await self.proxy_manager.release_proxy(proxy, release_success)
if use_socks5:
await self.socks5_pool.release(proxy)
if is_proxy_dead: if is_proxy_dead:
logger.warning(f"Visit attempt {attempt + 1} proxy connection failed, retrying with new proxy") logger.warning(f"Visit attempt {attempt + 1} proxy connection failed, retrying with new proxy")
continue continue
elif is_load_error: elif is_load_error:
# Ошибка загрузки — не меняем прокси, просто возвращаем результат
# (сайт мог быть временно недоступен, прокси ок)
logger.warning(f"Visit attempt {attempt + 1} page load error: {result.error or 'unknown'}") logger.warning(f"Visit attempt {attempt + 1} page load error: {result.error or 'unknown'}")
self._proxy_fail_streak = 0
return result return result
except (ValueError, InvalidIP) as e: except asyncio.TimeoutError:
logger.warning(f"Visit attempt {attempt + 1} timed out after {attempt_timeout}s, retrying")
except (ValueError, InvalidIP, InvalidProxy) as e:
logger.warning(f"Visit attempt {attempt + 1} proxy error, retrying: {e}") logger.warning(f"Visit attempt {attempt + 1} proxy error, retrying: {e}")
if proxy:
await self.proxy_manager.release_proxy(proxy, success=False)
if use_socks5:
await self.socks5_pool.release(proxy)
except Exception as e: except Exception as e:
logger.error(f"Visit error: {e}", exc_info=True) logger.error(f"Visit error: {e}", exc_info=True)
if proxy:
# Ошибка кода/браузера — прокси не виноват
await self.proxy_manager.release_proxy(proxy, success=True)
if use_socks5:
await self.socks5_pool.release(proxy)
return VisitResult(url=url, success=False, error=str(e)) return VisitResult(url=url, success=False, error=str(e))
self._proxy_fail_streak += 1
if self._proxy_fail_streak >= self._PROXY_FAIL_THRESHOLD:
logger.error(f"Proxy fail streak {self._proxy_fail_streak} — signalling PROXY_UNAVAILABLE")
self._proxy_fail_streak = 0
return VisitResult(url=url, success=False, error="PROXY_UNAVAILABLE")
return VisitResult(url=url, success=False, error="All attempts failed (fingerprint error)") return VisitResult(url=url, success=False, error="All attempts failed (fingerprint error)")
async def _browse_page( async def _browse_page(
@@ -326,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
@@ -373,13 +379,7 @@ class BrowserService:
return None return None
async def cleanup(self): async def cleanup(self):
"""Очистка с таймаутом.""" pass
try:
await asyncio.wait_for(self.socks5_pool.stop_all(), timeout=10)
except asyncio.TimeoutError:
logger.warning("SOCKS5 pool cleanup timeout")
except Exception as e:
logger.error(f"Cleanup error: {e}")
async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult: async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult:
return await self.visit_page(url, reading_time) return await self.visit_page(url, reading_time)
+275
View File
@@ -0,0 +1,275 @@
"""
Twitch чат через IRC over WebSocket (wss://irc-ws.chat.twitch.tv:443).
Замена raw TCP IRC — тот же протокол, WebSocket транспорт.
"""
import asyncio
import logging
import random
import re
from typing import Optional, Callable, Awaitable, List
from datetime import datetime
from urllib.parse import urlparse
import aiohttp
from config.settings import settings
logger = logging.getLogger(__name__)
_IRC_WS_URL = "wss://irc-ws.chat.twitch.tv:443"
class TwitchGQLChatClient:
"""
Twitch IRC over WebSocket клиент.
Интерфейс совместим с TwitchIRCClient.
"""
def __init__(
self,
channel: str,
target_username: str,
irc_username: str = None,
irc_oauth: str = None,
):
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
self._session: Optional[aiohttp.ClientSession] = None
self._connected = False
self._url_pattern = re.compile(r'https?://[^\s<>"]+')
async def connect(self) -> bool:
try:
self._session = aiohttp.ClientSession()
self._ws = await self._session.ws_connect(
_IRC_WS_URL,
timeout=aiohttp.ClientTimeout(total=15),
)
await self._send(f"PASS {self.irc_oauth}")
await self._send(f"NICK {self.irc_username}")
await self._send("CAP REQ :twitch.tv/tags")
await self._send("CAP REQ :twitch.tv/commands")
await self._send(f"JOIN #{self.channel}")
await asyncio.sleep(2)
self._connected = True
logger.info(f"✅ IRC WS подключен к #{self.channel}")
return True
except Exception as e:
logger.error(f"❌ IRC WS ошибка подключения: {e}")
await self._cleanup()
return False
async def _send(self, message: str) -> None:
if self._ws and not self._ws.closed:
await self._ws.send_str(f"{message}\r\n")
async def _cleanup(self) -> None:
self._connected = False
try:
if self._ws and not self._ws.closed:
await self._ws.close()
except Exception:
pass
try:
if self._session and not self._session.closed:
await self._session.close()
except Exception:
pass
self._ws = None
self._session = None
def _parse_message(self, line: str) -> Optional[dict]:
if not line or "PRIVMSG" not in line:
return None
try:
tags = {}
if line.startswith("@"):
tags_part = line.split(" ", 1)[0]
for tag in tags_part.lstrip("@").split(";"):
if "=" in tag:
k, v = tag.split("=", 1)
tags[k] = v
display_name = tags.get("display-name", "")
message_text = line.split(" :", 1)[1] if " :" in line else ""
if not message_text.strip():
return None
return {
"display_name": display_name or "unknown",
"message": message_text,
"tags": tags,
}
except Exception as e:
logger.error(f"Parse error: {e}")
return None
def _extract_urls(self, text: str) -> List[str]:
urls = []
for raw_url in self._url_pattern.findall(text):
try:
parsed = urlparse(raw_url)
clean = f"{parsed.scheme}://{parsed.netloc}{parsed.path}"
if parsed.query:
clean += f"?{parsed.query}"
urls.append(clean)
except Exception:
urls.append(raw_url.strip())
return urls
def _is_domain_allowed(self, url: str, allowed_domains: List[str]) -> bool:
if not allowed_domains:
return True
try:
domain = urlparse(url).netloc.lower()
return any(d.lower() in domain for d in allowed_domains)
except Exception:
return False
async def listen_for_messages(
self,
on_url_found: Callable[[str, str], Awaitable[None]],
duration: int,
allowed_domains: Optional[List[str]] = None,
is_active: Callable[[], bool] = None,
on_drain_fail: Callable[[], Awaitable[None]] = None,
) -> dict:
stats = {"links_found": 0, "messages": 0, "reconnects": 0}
last_message_time = datetime.now()
drain_errors = 0
if is_active and not is_active():
logger.info(f"⏸️ #{self.channel} ждём resume перед подключением...")
while is_active and not is_active():
await asyncio.sleep(2)
if not await self.connect():
raise ConnectionError(f"IRC WS connect failed for #{self.channel}")
# start_time после паузы и подключения — пауза не считается в длительность
start_time = datetime.now()
while duration == 0 or (datetime.now() - start_time).total_seconds() < duration:
if is_active and not is_active():
logger.info(f"⏸️ #{self.channel} пауза...")
await self._cleanup()
pause_begin = datetime.now()
while is_active and not is_active():
await asyncio.sleep(2)
# Исключаем время паузы из отсчёта длительности
start_time += datetime.now() - pause_begin
logger.info(f"▶️ #{self.channel} resume, переподключение...")
if not await self.connect():
drain_errors += 1
if drain_errors >= 3 and on_drain_fail:
await on_drain_fail()
return stats
else:
drain_errors = 0
stats["reconnects"] += 1
last_message_time = datetime.now()
continue
if not self._connected or self._ws is None or self._ws.closed:
await asyncio.sleep(5)
if not await self.connect():
drain_errors += 1
if drain_errors >= 3 and on_drain_fail:
await on_drain_fail()
return stats
else:
drain_errors = 0
stats["reconnects"] += 1
last_message_time = datetime.now()
continue
# Реконнект при отсутствии активности > 6 минут
if (datetime.now() - last_message_time).total_seconds() > 360:
logger.info("🔄 Нет активности 6 мин, переподключение...")
await self._cleanup()
await asyncio.sleep(1)
if not await self.connect():
drain_errors += 1
if drain_errors >= 3 and on_drain_fail:
await on_drain_fail()
return stats
else:
drain_errors = 0
stats["reconnects"] += 1
last_message_time = datetime.now()
continue
try:
raw = await asyncio.wait_for(self._ws.receive(), timeout=5.0)
except asyncio.CancelledError:
raise
except asyncio.TimeoutError:
continue
except Exception as e:
logger.warning(f"IRC WS receive error: {e}")
self._connected = False
continue
if raw.type in (
aiohttp.WSMsgType.CLOSING,
aiohttp.WSMsgType.CLOSED,
aiohttp.WSMsgType.ERROR,
):
logger.warning(f"WS закрыт: {raw.type}")
self._connected = False
continue
if raw.type != aiohttp.WSMsgType.TEXT:
continue
line = raw.data.strip()
if not line:
continue
if line.startswith("PING"):
await self._send("PONG :tmi.twitch.tv")
last_message_time = datetime.now()
logger.debug("🏓 PONG")
continue
last_message_time = datetime.now()
msg = self._parse_message(line)
if not msg:
continue
stats["messages"] += 1
logger.debug(f"IRC WS msg from {msg['display_name']}: {msg['message'][:50]}...")
for url in self._extract_urls(msg["message"]):
if self._is_domain_allowed(url, allowed_domains):
stats["links_found"] += 1
try:
await on_url_found(url, msg["display_name"])
except asyncio.CancelledError:
raise
except Exception:
pass
else:
logger.info(
f"🚫 Домен не разрешён: {url} (allowed_domains={allowed_domains})"
)
return stats
async def disconnect(self) -> None:
await self._cleanup()
+15 -5
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
@@ -173,7 +180,6 @@ class TwitchIRCClient:
) -> dict: ) -> dict:
stats = {"links_found": 0, "messages": 0, "reconnects": 0} stats = {"links_found": 0, "messages": 0, "reconnects": 0}
start_time = datetime.now()
last_message_time = datetime.now() last_message_time = datetime.now()
drain_errors = 0 drain_errors = 0
@@ -182,20 +188,24 @@ class TwitchIRCClient:
logger.info(f"⏸️ #{self.channel} waiting for resume before connecting...") logger.info(f"⏸️ #{self.channel} waiting for resume before connecting...")
while is_active and not is_active(): while is_active and not is_active():
await asyncio.sleep(2) await asyncio.sleep(2)
if not is_active and not is_active():
return stats
logger.info(f"▶️ #{self.channel} resumed, connecting to IRC") logger.info(f"▶️ #{self.channel} resumed, connecting to IRC")
if not await self.connect(): if not await self.connect():
raise ConnectionError(f"IRC connect failed for #{self.channel}") raise ConnectionError(f"IRC connect failed for #{self.channel}")
while duration == 0 or (datetime.now() - start_time).seconds < duration: # start_time после паузы и подключения — пауза не считается в длительность
start_time = datetime.now()
while duration == 0 or (datetime.now() - start_time).total_seconds() < duration:
# === ПРОВЕРКА АКТИВНОСТИ === # === ПРОВЕРКА АКТИВНОСТИ ===
if is_active and not is_active(): if is_active and not is_active():
logger.info(f"⏸️ #{self.channel} paused, waiting for resume...") logger.info(f"⏸️ #{self.channel} paused, waiting for resume...")
await self.disconnect() await self.disconnect()
pause_begin = datetime.now()
while is_active and not is_active(): while is_active and not is_active():
await asyncio.sleep(2) await asyncio.sleep(2)
# Исключаем время паузы из отсчёта длительности
start_time += datetime.now() - pause_begin
logger.info(f"▶️ #{self.channel} resumed, reconnecting") logger.info(f"▶️ #{self.channel} resumed, reconnecting")
await self.connect() await self.connect()
last_message_time = datetime.now() last_message_time = datetime.now()
+26 -17
View File
@@ -51,59 +51,68 @@ class VisitScheduler:
async def run(self) -> dict: async def run(self) -> dict:
""" """
Run scheduled visits. Run scheduled visits.
Returns: Returns:
Statistics dict Statistics dict
""" """
total_reading_time = 0 total_reading_time = 0
loop = asyncio.get_event_loop()
try: try:
while True: while True:
if self.max_visits and self.visit_count >= self.max_visits: if self.max_visits and self.visit_count >= self.max_visits:
logger.info(f"Visit limit reached: {self.max_visits}") logger.info(f"Visit limit reached: {self.max_visits}")
break break
# Randomize parameters # Randomize parameters once per cycle
current_reading = random.randint(self.min_reading, self.max_reading) current_reading = random.randint(self.min_reading, self.max_reading)
current_delay = random.randint(self.min_delay, self.max_delay) current_delay = random.randint(self.min_delay, self.max_delay)
self.visit_count += 1 self.visit_count += 1
# Record start time so the next visit fires strictly on the timer
tick_start = loop.time()
logger.info( logger.info(
f"Visit {self.visit_count}/{'' if not self.max_visits else self.max_visits}: " f"Visit {self.visit_count}/{'' if not self.max_visits else self.max_visits}: "
f"reading={current_reading}s" f"reading={current_reading}s delay={current_delay}s"
) )
# Visit page # Visit page
result = await self.browser_service.visit_page( result = await self.browser_service.visit_page(
self.url, current_reading self.url, current_reading
) )
if result.success: if result.success:
self.successful += 1 self.successful += 1
total_reading_time += current_reading total_reading_time += current_reading
else: else:
self.failed += 1 self.failed += 1
# Notify about visit completion # Notify about visit completion
if self.on_visit_complete: if self.on_visit_complete:
await self.on_visit_complete( await self.on_visit_complete(
self.visit_count, result, current_delay, current_reading self.visit_count, result, current_delay, current_reading
) )
if self.max_visits and self.visit_count >= self.max_visits: if self.max_visits and self.visit_count >= self.max_visits:
break break
# Wait before next visit # Wait only the remaining time so that the period between visit
# starts is exactly current_delay (independent of visit duration)
elapsed = loop.time() - tick_start
remaining = current_delay - elapsed
if self.on_progress: if self.on_progress:
await self.on_progress(self.visit_count, current_delay) await self.on_progress(self.visit_count, current_delay)
await asyncio.sleep(current_delay) if remaining > 0:
await asyncio.sleep(remaining)
except asyncio.CancelledError: except asyncio.CancelledError:
logger.info("Visit scheduler cancelled") logger.info("Visit scheduler cancelled")
raise raise
return { return {
"total_visits": self.visit_count, "total_visits": self.visit_count,
"successful": self.successful, "successful": self.successful,
+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:
+336
View File
@@ -0,0 +1,336 @@
"""
Утилита проверки прокси.
- Удаляет дубликаты по ip:port
- Тестирует напрямую (HTTP и SOCKS5)
- Тестирует через локальный SOCKS5→HTTP туннель
- Сохраняет результат обратно в файл
Использование:
python tools/proxy_checker.py
python tools/proxy_checker.py --file proxies.txt --concurrency 10
python tools/proxy_checker.py --no-save # только показать, не писать файл
"""
import asyncio
import sys
import time
import argparse
from pathlib import Path
from dataclasses import dataclass, field
from typing import Optional
sys.path.insert(0, str(Path(__file__).parent.parent))
from managers.proxy_manager import Proxy, ProxyType
TEST_URL = "https://www.google.com/generate_204"
TIMEOUT = 10.0
@dataclass
class CheckResult:
proxy: Proxy
direct_ok: bool = False
direct_ms: Optional[float] = None
tunnel_ok: bool = False
tunnel_ms: Optional[float] = None
error: Optional[str] = None
@property
def ok(self) -> bool:
return self.direct_ok or self.tunnel_ok
def summary(self) -> str:
parts = []
if self.direct_ok:
parts.append(f"direct {self.direct_ms:.0f}ms")
if self.tunnel_ok:
parts.append(f"tunnel {self.tunnel_ms:.0f}ms")
if not parts:
return f"FAIL {self.error or ''}"
return "OK " + " | ".join(parts)
async def _check_direct(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Прямая проверка через requests в thread pool."""
import asyncio
loop = asyncio.get_running_loop()
def _sync():
import requests
t0 = time.perf_counter()
scheme = proxy.proxy_type.value
if proxy.login and proxy.password:
url = f"{scheme}://{proxy.login}:{proxy.password}@{proxy.ip}:{proxy.port}"
else:
url = f"{scheme}://{proxy.ip}:{proxy.port}"
proxies = {"http": url, "https": url}
try:
r = requests.get(TEST_URL, proxies=proxies, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000
return r.status_code in (200, 204), ms
except Exception:
return False, None
return await loop.run_in_executor(None, _sync)
async def _check_tunnel(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Проверка через SOCKS5→HTTP туннель (только для SOCKS5/SOCKS4)."""
if proxy.proxy_type not in (ProxyType.SOCKS5, ProxyType.SOCKS4):
return False, None
from services.socks5_to_http_proxy import Socks5ToHttpProxy
import aiohttp
tunnel = Socks5ToHttpProxy(
socks5_host=proxy.ip,
socks5_port=int(proxy.port),
username=proxy.login,
password=proxy.password,
)
try:
await asyncio.wait_for(tunnel.start(), timeout=5)
t0 = time.perf_counter()
ok = await tunnel.check_connection(test_url=TEST_URL, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000 if ok else None
return ok, ms
except Exception:
return False, None
finally:
try:
await asyncio.wait_for(tunnel.stop(), timeout=3)
except Exception:
pass
async def _check_direct_as_http(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Принудительная проверка как HTTP прокси (игнорирует настроенный тип)."""
loop = asyncio.get_running_loop()
def _sync():
import requests
t0 = time.perf_counter()
if proxy.login and proxy.password:
url = f"http://{proxy.login}:{proxy.password}@{proxy.ip}:{proxy.port}"
else:
url = f"http://{proxy.ip}:{proxy.port}"
proxies = {"http": url, "https": url}
try:
r = requests.get(TEST_URL, proxies=proxies, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000
return r.status_code in (200, 204), ms
except Exception:
return False, None
return await loop.run_in_executor(None, _sync)
async def check_proxy(proxy: Proxy, timeout: float) -> CheckResult:
result = CheckResult(proxy=proxy)
try:
direct_ok, direct_ms = await asyncio.wait_for(
_check_direct(proxy, timeout), timeout=timeout + 2
)
result.direct_ok = direct_ok
result.direct_ms = direct_ms
if proxy.proxy_type in (ProxyType.SOCKS5, ProxyType.SOCKS4):
tunnel_ok, tunnel_ms = await asyncio.wait_for(
_check_tunnel(proxy, timeout), timeout=timeout + 5
)
result.tunnel_ok = tunnel_ok
result.tunnel_ms = tunnel_ms
except Exception as e:
result.error = str(e)[:60]
return result
async def retry_failed_as_http(results: list[CheckResult], timeout: float, sem: asyncio.Semaphore) -> dict[str, CheckResult]:
"""
Повторно проверяет упавшие прокси принудительно как HTTP.
Возвращает dict[ip:port -> CheckResult] только для тех что теперь прошли.
"""
failed = [r for r in results if not r.ok]
if not failed:
return {}
recovered: dict[str, CheckResult] = {}
async def _retry(r: CheckResult):
async with sem:
try:
ok, ms = await asyncio.wait_for(
_check_direct_as_http(r.proxy, timeout), timeout=timeout + 2
)
except (TimeoutError, asyncio.TimeoutError, asyncio.CancelledError, Exception):
ok, ms = False, None
if ok:
key = f"{r.proxy.ip}:{r.proxy.port}"
r.direct_ok = True
r.direct_ms = ms
r.error = None
r.proxy.proxy_type = ProxyType.HTTP
recovered[key] = r
await asyncio.gather(*[_retry(r) for r in failed], return_exceptions=True)
return recovered
def _load_proxies(file_path: str) -> tuple[list[Proxy], list[str]]:
"""Загружает прокси из файла. Возвращает (список прокси, сырые строки)."""
raw_lines = Path(file_path).read_text(encoding="utf-8").splitlines()
proxies = []
seen = set() # ip:port для дедупликации
for line in raw_lines:
stripped = line.strip()
if not stripped or stripped.startswith("#") or stripped.startswith("!"):
continue
p = Proxy.from_line(stripped)
if not p:
continue
key = f"{p.ip}:{p.port}"
if key in seen:
continue
seen.add(key)
proxies.append(p)
return proxies, raw_lines
def _save_results(file_path: str, results: list[CheckResult], raw_lines: list[str]):
"""Перезаписывает файл: рабочие без '!', нерабочие с '!'.
Использует актуальный тип прокси из results (может быть изменён retry-фазой).
"""
# ip:port → актуальный объект Proxy с возможно обновлённым типом
proxy_map: dict[str, Proxy] = {
f"{r.proxy.ip}:{r.proxy.port}": r.proxy for r in results
}
working = {f"{r.proxy.ip}:{r.proxy.port}" for r in results if r.ok}
new_lines = []
seen_keys: set[str] = set() # для удаления дубликатов
for line in raw_lines:
stripped = line.strip()
if not stripped or stripped.startswith("#"):
new_lines.append(line)
continue
clean = stripped.lstrip("!")
p = Proxy.from_line(clean)
if not p:
new_lines.append(line)
continue
key = f"{p.ip}:{p.port}"
if key in seen_keys:
continue # дубликат — пропускаем
seen_keys.add(key)
actual = proxy_map.get(key, p)
scheme = actual.proxy_type.value
if actual.login and actual.password:
formatted = f"{scheme}://{actual.ip}:{actual.port}:{actual.login}:{actual.password}"
else:
formatted = f"{scheme}://{actual.ip}:{actual.port}"
if key in working:
new_lines.append(formatted)
else:
new_lines.append(f"!{formatted}")
Path(file_path).write_text("\n".join(new_lines) + "\n", encoding="utf-8")
async def run_checker(file_path: str, concurrency: int, timeout: float, save: bool, retry_http: bool):
proxies, raw_lines = _load_proxies(file_path)
total_raw = sum(
1 for l in raw_lines
if l.strip() and not l.strip().startswith("#") and not l.strip().startswith("!")
)
duplicates = total_raw - len(proxies)
print(f"\n{'='*60}")
print(f" Файл: {file_path}")
print(f" Загружено: {total_raw} | Дубликатов: {duplicates} | К проверке: {len(proxies)}")
print(f" Параллельность: {concurrency} | Таймаут: {timeout}s")
if retry_http:
print(f" Повтор упавших через HTTP: ВКЛ")
print(f"{'='*60}\n")
sem = asyncio.Semaphore(concurrency)
results: list[CheckResult] = []
done = 0
async def _run(proxy: Proxy):
nonlocal done
async with sem:
r = await check_proxy(proxy, timeout)
results.append(r)
done += 1
status = "" if r.ok else ""
print(f" [{done:>3}/{len(proxies)}] {status} {proxy.id:<35} {r.summary()}")
await asyncio.gather(*[_run(p) for p in proxies])
# Фаза 2: повторная проверка упавших как HTTP
recovered_count = 0
if retry_http:
failed_count = sum(1 for r in results if not r.ok)
if failed_count:
print(f"\n ── Повтор {failed_count} упавших как HTTP прямую ──\n")
recovered = await retry_failed_as_http(results, timeout, sem)
recovered_count = len(recovered)
for key, r in recovered.items():
print(f" [retry] ✅ {r.proxy.ip}:{r.proxy.port} → HTTP {r.direct_ms:.0f}ms")
ok_count = sum(1 for r in results if r.ok)
fail_count = len(results) - ok_count
direct_ok = sum(1 for r in results if r.direct_ok)
tunnel_ok = sum(1 for r in results if r.tunnel_ok)
print(f"\n{'='*60}")
print(f" ИТОГО")
print(f"{'='*60}")
print(f" Рабочих : {ok_count} / {len(proxies)}")
print(f" Нерабочих : {fail_count}")
print(f" Direct OK : {direct_ok}")
print(f" Tunnel OK : {tunnel_ok}")
if recovered_count:
print(f" Спасено HTTP: {recovered_count}")
if duplicates:
print(f" Дубликатов : {duplicates} (удалены)")
if save:
_save_results(file_path, results, raw_lines)
print(f"\n Файл обновлён: {file_path}")
print(f" Нерабочие помечены '!'")
print(f"{'='*60}\n")
def main():
parser = argparse.ArgumentParser(description="Проверка прокси")
parser.add_argument("--file", default="proxies.txt", help="Файл с прокси")
parser.add_argument("--concurrency", type=int, default=5, help="Параллельность (по умолчанию: 5)")
parser.add_argument("--timeout", type=float, default=10.0, help="Таймаут секунд (по умолчанию: 10)")
parser.add_argument("--no-save", action="store_true", help="Не перезаписывать файл")
parser.add_argument("--retry-http", action="store_true", help="Повторить упавшие как HTTP напрямую")
args = parser.parse_args()
asyncio.run(run_checker(
file_path=args.file,
concurrency=args.concurrency,
timeout=args.timeout,
save=not args.no_save,
retry_http=args.retry_http,
))
if __name__ == "__main__":
main()