Compare commits
10 Commits
22b49109fa
...
f262616aad
| Author | SHA1 | Date | |
|---|---|---|---|
| f262616aad | |||
| e438f09693 | |||
| b02b7518cf | |||
| 61ff90e51b | |||
| 493c67cb1d | |||
| 795b2e928c | |||
| d95dc462d9 | |||
| 406a58c585 | |||
| 68275386bf | |||
| e2e2761b51 |
@@ -1,8 +1,10 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
import hashlib
|
import hashlib
|
||||||
|
import json
|
||||||
import secrets
|
import secrets
|
||||||
import logging
|
import logging
|
||||||
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,6 +12,7 @@ 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:
|
||||||
@@ -22,6 +25,7 @@ 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 _hash_password(password: str, salt: str = None) -> tuple:
|
||||||
@@ -37,6 +41,10 @@ class AuthManager:
|
|||||||
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 +57,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
|
||||||
|
|||||||
+6
-5
@@ -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 = ""
|
||||||
|
|
||||||
|
|||||||
+152
-134
@@ -22,7 +22,6 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder, ReplyKeyboardBuilder
|
|||||||
|
|
||||||
from config.settings import settings
|
from config.settings import settings
|
||||||
from managers.background_tasks import BackgroundTaskManager
|
from managers.background_tasks import BackgroundTaskManager
|
||||||
from managers.proxy_manager import ProxyManager
|
|
||||||
from managers.task_manager import TaskManager, TaskParams
|
from managers.task_manager import TaskManager, TaskParams
|
||||||
from services.browser_service import BrowserService
|
from services.browser_service import BrowserService
|
||||||
from services.browser_pool import BrowserPool
|
from services.browser_pool import BrowserPool
|
||||||
@@ -45,7 +44,6 @@ class BotInterface:
|
|||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
background_tasks: BackgroundTaskManager,
|
background_tasks: BackgroundTaskManager,
|
||||||
proxy_manager: ProxyManager,
|
|
||||||
browser_service: BrowserService,
|
browser_service: BrowserService,
|
||||||
browser_pool: BrowserPool = None,
|
browser_pool: BrowserPool = None,
|
||||||
storage: TaskStorage = None,
|
storage: TaskStorage = None,
|
||||||
@@ -53,7 +51,6 @@ class BotInterface:
|
|||||||
bot_ref=None,
|
bot_ref=None,
|
||||||
):
|
):
|
||||||
self.background_tasks = background_tasks
|
self.background_tasks = background_tasks
|
||||||
self.proxy_manager = proxy_manager
|
|
||||||
self.browser_service = browser_service
|
self.browser_service = browser_service
|
||||||
self.browser_pool = browser_pool
|
self.browser_pool = browser_pool
|
||||||
self.task_manager = TaskManager()
|
self.task_manager = TaskManager()
|
||||||
@@ -72,6 +69,7 @@ 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_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] = {}
|
||||||
@@ -258,29 +256,6 @@ class BotInterface:
|
|||||||
await interface._show_status(callback.message, edit=True)
|
await interface._show_status(callback.message, edit=True)
|
||||||
await callback.answer()
|
await callback.answer()
|
||||||
|
|
||||||
# === Перезапуск ===
|
|
||||||
@dp.callback_query(F.data == "reload_proxies")
|
|
||||||
async def cb_reload_proxies(callback: CallbackQuery):
|
|
||||||
if not await require_admin(callback):
|
|
||||||
return
|
|
||||||
count = interface.proxy_manager.reload()
|
|
||||||
await interface._show_status(callback.message, edit=True)
|
|
||||||
await callback.answer(f"✅ Прокси перезагружены: {count} шт.")
|
|
||||||
|
|
||||||
@dp.callback_query(F.data == "check_proxies")
|
|
||||||
async def cb_check_proxies(callback: CallbackQuery):
|
|
||||||
if not await require_admin(callback):
|
|
||||||
return
|
|
||||||
total = interface.proxy_manager.count
|
|
||||||
if total == 0:
|
|
||||||
await callback.answer("❌ Нет загруженных прокси", show_alert=True)
|
|
||||||
return
|
|
||||||
await callback.message.edit_text(
|
|
||||||
f"⏳ Проверка {total} прокси через туннель...\n\nЭто займёт несколько секунд."
|
|
||||||
)
|
|
||||||
await callback.answer()
|
|
||||||
asyncio.create_task(interface._check_all_proxies(callback.message))
|
|
||||||
|
|
||||||
@dp.callback_query(F.data == "bot_restart_confirm")
|
@dp.callback_query(F.data == "bot_restart_confirm")
|
||||||
async def cb_restart_confirm(callback: CallbackQuery):
|
async def cb_restart_confirm(callback: CallbackQuery):
|
||||||
if not await require_admin(callback):
|
if not await require_admin(callback):
|
||||||
@@ -508,9 +483,14 @@ class BotInterface:
|
|||||||
if not await require_admin(callback):
|
if not await require_admin(callback):
|
||||||
return
|
return
|
||||||
tid = callback.data.replace("tdelete_", "", 1)
|
tid = callback.data.replace("tdelete_", "", 1)
|
||||||
|
params = await interface.task_manager.get_task(tid)
|
||||||
|
is_visit = params and params.task_type in ("visit", "user_visit")
|
||||||
await interface.task_manager.remove_task(tid)
|
await interface.task_manager.remove_task(tid)
|
||||||
await callback.answer("🗑️ Удалена")
|
await callback.answer("🗑️ Удалена")
|
||||||
await interface._show_streamers_list(callback.message, edit=True)
|
if is_visit:
|
||||||
|
await interface._show_tasks_list(callback.message, edit=True)
|
||||||
|
else:
|
||||||
|
await interface._show_streamers_list(callback.message, edit=True)
|
||||||
|
|
||||||
|
|
||||||
@dp.callback_query(F.data.startswith("treset_"))
|
@dp.callback_query(F.data.startswith("treset_"))
|
||||||
@@ -1239,6 +1219,28 @@ class BotInterface:
|
|||||||
await interface._show_user_detail(callback, target_uid)
|
await interface._show_user_detail(callback, target_uid)
|
||||||
await callback.answer()
|
await callback.answer()
|
||||||
|
|
||||||
|
@dp.callback_query(F.data.startswith("umsg_"))
|
||||||
|
async def cb_admin_msg(callback: CallbackQuery):
|
||||||
|
if not await require_admin(callback):
|
||||||
|
return
|
||||||
|
target_uid = int(callback.data.replace("umsg_", "", 1))
|
||||||
|
interface._admin_msg_state[callback.from_user.id] = target_uid
|
||||||
|
builder = InlineKeyboardBuilder()
|
||||||
|
builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data=f"udetail_{target_uid}"))
|
||||||
|
await callback.message.edit_text(
|
||||||
|
f"✉️ Сообщение пользователю {target_uid}\n\nВведите текст:",
|
||||||
|
reply_markup=builder.as_markup()
|
||||||
|
)
|
||||||
|
await callback.answer()
|
||||||
|
|
||||||
|
@dp.callback_query(F.data == "dismiss_msg")
|
||||||
|
async def cb_dismiss_msg(callback: CallbackQuery):
|
||||||
|
try:
|
||||||
|
await callback.message.delete()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
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):
|
||||||
@@ -1274,6 +1276,25 @@ class BotInterface:
|
|||||||
user_id = message.from_user.id
|
user_id = message.from_user.id
|
||||||
text = message.text.strip() if message.text else ""
|
text = message.text.strip() if message.text else ""
|
||||||
|
|
||||||
|
# Админ пишет сообщение пользователю
|
||||||
|
if user_id in interface._admin_msg_state:
|
||||||
|
if not text:
|
||||||
|
return
|
||||||
|
target_uid = interface._admin_msg_state.pop(user_id)
|
||||||
|
await interface._safe_delete(message.bot, message.chat.id, message.message_id)
|
||||||
|
builder = InlineKeyboardBuilder()
|
||||||
|
builder.row(InlineKeyboardButton(text="✅ Понял", callback_data="dismiss_msg"))
|
||||||
|
try:
|
||||||
|
await message.bot.send_message(
|
||||||
|
target_uid,
|
||||||
|
f"📢 Сообщение от администратора:\n\n{text}",
|
||||||
|
reply_markup=builder.as_markup()
|
||||||
|
)
|
||||||
|
await message.answer(f"✅ Сообщение отправлено пользователю {target_uid}")
|
||||||
|
except Exception as e:
|
||||||
|
await message.answer(f"❌ Не удалось отправить: {e}")
|
||||||
|
return
|
||||||
|
|
||||||
# Ввод пароля (неавторизованный пользователь ждёт пароль)
|
# Ввод пароля (неавторизованный пользователь ждёт пароль)
|
||||||
if user_id in interface._waiting_password:
|
if user_id in interface._waiting_password:
|
||||||
if not text:
|
if not text:
|
||||||
@@ -1674,9 +1695,67 @@ class BotInterface:
|
|||||||
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):
|
||||||
"""Все задачи."""
|
"""Все задачи — посещения и мониторинг."""
|
||||||
text = await self.task_manager.format_task_list()
|
active = await self.task_manager.get_active_tasks()
|
||||||
|
completed = await self.task_manager.get_completed_tasks()
|
||||||
|
|
||||||
|
visit_active = {tid: p for tid, p in active.items() if p.task_type in ("visit", "user_visit")}
|
||||||
|
visit_done = {tid: p for tid, p in completed.items() if p.task_type in ("visit", "user_visit")}
|
||||||
|
twitch_active = {tid: p for tid, p in active.items() if p.task_type == "twitch_irc"}
|
||||||
|
twitch_done = {tid: p for tid, p in completed.items() if p.task_type == "twitch_irc"}
|
||||||
|
|
||||||
builder = InlineKeyboardBuilder()
|
builder = InlineKeyboardBuilder()
|
||||||
|
|
||||||
|
if not any([visit_active, visit_done, twitch_active, twitch_done]):
|
||||||
|
text = "📊 Задачи\n\nЗадач нет."
|
||||||
|
else:
|
||||||
|
text = "📊 Задачи\n\n"
|
||||||
|
|
||||||
|
if twitch_active:
|
||||||
|
text += "📺 Twitch — активные:\n"
|
||||||
|
for tid, p in list(twitch_active.items())[:10]:
|
||||||
|
em = p.get_status_emoji()
|
||||||
|
domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все"
|
||||||
|
text += f"{em} {p.channel}\n 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n"
|
||||||
|
builder.row(InlineKeyboardButton(
|
||||||
|
text=f"{em} 📺 {p.channel} (🔗{p.links_found})",
|
||||||
|
callback_data=f"tdetail_{tid}"
|
||||||
|
))
|
||||||
|
|
||||||
|
if twitch_done:
|
||||||
|
text += "\n📺 Twitch — завершённые:\n"
|
||||||
|
for tid, p in list(twitch_done.items())[:5]:
|
||||||
|
em = p.get_status_emoji()
|
||||||
|
text += f"{em} {p.channel}\n"
|
||||||
|
builder.row(InlineKeyboardButton(
|
||||||
|
text=f"{em} 📺 {p.channel} (завершена)",
|
||||||
|
callback_data=f"tdetail_{tid}"
|
||||||
|
))
|
||||||
|
|
||||||
|
if visit_active:
|
||||||
|
text += "\n🔗 Посещения — активные:\n"
|
||||||
|
for tid, p in list(visit_active.items())[:10]:
|
||||||
|
em = p.get_status_emoji()
|
||||||
|
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 ✅{p.successful_visits}/{p.max_visits} | 📊{p.total_visits}\n"
|
||||||
|
builder.row(InlineKeyboardButton(
|
||||||
|
text=f"{em} 🔗 {url_s[:28]} ({p.successful_visits}/{p.max_visits})",
|
||||||
|
callback_data=f"tdetail_{tid}"
|
||||||
|
))
|
||||||
|
|
||||||
|
if visit_done:
|
||||||
|
text += "\n🔗 Посещения — завершённые:\n"
|
||||||
|
for tid, p in list(visit_done.items())[:5]:
|
||||||
|
em = p.get_status_emoji()
|
||||||
|
url_s = (p.url or "")[:35]
|
||||||
|
uid_tag = f" | uid:{p.user_id}" if p.user_id else ""
|
||||||
|
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"))
|
||||||
await self._edit_or_send(message, text, builder.as_markup(), edit)
|
await self._edit_or_send(message, text, builder.as_markup(), edit)
|
||||||
@@ -1700,7 +1779,6 @@ class BotInterface:
|
|||||||
|
|
||||||
async def _show_status(self, message: Message, edit: bool = False):
|
async def _show_status(self, message: Message, edit: bool = False):
|
||||||
stats = await self.task_manager.get_stats()
|
stats = await self.task_manager.get_stats()
|
||||||
proxy_stats = self.proxy_manager.get_stats()
|
|
||||||
bg_count = self.background_tasks.active_count
|
bg_count = self.background_tasks.active_count
|
||||||
|
|
||||||
delta = datetime.now() - self._start_time
|
delta = datetime.now() - self._start_time
|
||||||
@@ -1718,14 +1796,11 @@ class BotInterface:
|
|||||||
f"├ 🌐 Визиты: {stats['visit']}\n"
|
f"├ 🌐 Визиты: {stats['visit']}\n"
|
||||||
f"├ ⏸️ Пауза: {stats['paused']}\n"
|
f"├ ⏸️ Пауза: {stats['paused']}\n"
|
||||||
f"└ 🔄 Активно: {stats['active']}\n\n"
|
f"└ 🔄 Активно: {stats['active']}\n\n"
|
||||||
f"🔌 Прокси: {proxy_stats['total']} / {proxy_stats['available']} доступно\n"
|
|
||||||
f"🔁 Фоновых задач: {bg_count}\n\n"
|
f"🔁 Фоновых задач: {bg_count}\n\n"
|
||||||
f"{sys_info}"
|
f"{sys_info}"
|
||||||
)
|
)
|
||||||
builder = InlineKeyboardBuilder()
|
builder = InlineKeyboardBuilder()
|
||||||
builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_status"))
|
builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_status"))
|
||||||
builder.row(InlineKeyboardButton(text="🔃 Перезагрузить прокси", callback_data="reload_proxies"))
|
|
||||||
builder.row(InlineKeyboardButton(text="🔍 Проверить прокси", callback_data="check_proxies"))
|
|
||||||
builder.row(InlineKeyboardButton(text="🔁 Перезапустить бота", callback_data="bot_restart_confirm"))
|
builder.row(InlineKeyboardButton(text="🔁 Перезапустить бота", callback_data="bot_restart_confirm"))
|
||||||
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)
|
||||||
@@ -1743,89 +1818,6 @@ class BotInterface:
|
|||||||
await asyncio.sleep(0.5)
|
await asyncio.sleep(0.5)
|
||||||
os.execv(sys.executable, [sys.executable] + sys.argv)
|
os.execv(sys.executable, [sys.executable] + sys.argv)
|
||||||
|
|
||||||
async def _check_all_proxies(self, message) -> None:
|
|
||||||
"""Проверяет все загруженные прокси через туннель и редактирует сообщение с результатами."""
|
|
||||||
from services.socks5_to_http_proxy import Socks5ToHttpProxy
|
|
||||||
from managers.proxy_manager import ProxyType
|
|
||||||
|
|
||||||
proxies = list(self.proxy_manager._proxies)
|
|
||||||
if not proxies:
|
|
||||||
builder = InlineKeyboardBuilder()
|
|
||||||
builder.row(InlineKeyboardButton(text="🔙 К статусу", callback_data="menu_status"))
|
|
||||||
try:
|
|
||||||
await message.edit_text("❌ Нет загруженных прокси", reply_markup=builder.as_markup())
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
return
|
|
||||||
|
|
||||||
sem = asyncio.Semaphore(5)
|
|
||||||
results: dict[str, bool] = {}
|
|
||||||
|
|
||||||
async def _check_one(proxy):
|
|
||||||
async with sem:
|
|
||||||
try:
|
|
||||||
if proxy.proxy_type == ProxyType.SOCKS5:
|
|
||||||
tunnel = Socks5ToHttpProxy(
|
|
||||||
socks5_host=proxy.ip,
|
|
||||||
socks5_port=int(proxy.port),
|
|
||||||
username=proxy.login,
|
|
||||||
password=proxy.password,
|
|
||||||
)
|
|
||||||
await tunnel.start()
|
|
||||||
try:
|
|
||||||
ok = await tunnel.check_connection(timeout=8.0)
|
|
||||||
finally:
|
|
||||||
await tunnel.stop()
|
|
||||||
else:
|
|
||||||
from aiohttp import ClientSession, ClientTimeout
|
|
||||||
proxy_url = proxy.server
|
|
||||||
if proxy.login and proxy.password:
|
|
||||||
from urllib.parse import urlparse
|
|
||||||
parsed = urlparse(proxy_url)
|
|
||||||
proxy_url = f"{parsed.scheme}://{proxy.login}:{proxy.password}@{parsed.netloc}"
|
|
||||||
async with ClientSession(timeout=ClientTimeout(total=8)) as s:
|
|
||||||
async with s.get(
|
|
||||||
"https://www.google.com/generate_204",
|
|
||||||
proxy=proxy_url,
|
|
||||||
allow_redirects=False,
|
|
||||||
) as resp:
|
|
||||||
ok = resp.status in (200, 204)
|
|
||||||
except Exception:
|
|
||||||
ok = False
|
|
||||||
results[proxy.id] = ok
|
|
||||||
|
|
||||||
await asyncio.gather(*[_check_one(p) for p in proxies], return_exceptions=True)
|
|
||||||
|
|
||||||
await self.proxy_manager.apply_check_results(results)
|
|
||||||
|
|
||||||
ok_ids = [pid for pid, ok in results.items() if ok]
|
|
||||||
fail_ids = [pid for pid, ok in results.items() if not ok]
|
|
||||||
|
|
||||||
lines = [
|
|
||||||
f"🔍 Проверка прокси завершена\n",
|
|
||||||
f"✅ Рабочих: {len(ok_ids)} / {len(proxies)}",
|
|
||||||
f"❌ Нерабочих: {len(fail_ids)}\n",
|
|
||||||
]
|
|
||||||
if ok_ids:
|
|
||||||
lines.append("✅ Работают:")
|
|
||||||
for pid in ok_ids[:15]:
|
|
||||||
lines.append(f" • {pid}")
|
|
||||||
if len(ok_ids) > 15:
|
|
||||||
lines.append(f" ...ещё {len(ok_ids) - 15}")
|
|
||||||
if fail_ids:
|
|
||||||
lines.append("\n❌ Не работают:")
|
|
||||||
for pid in fail_ids[:15]:
|
|
||||||
lines.append(f" • {pid}")
|
|
||||||
if len(fail_ids) > 15:
|
|
||||||
lines.append(f" ...ещё {len(fail_ids) - 15}")
|
|
||||||
|
|
||||||
builder = InlineKeyboardBuilder()
|
|
||||||
builder.row(InlineKeyboardButton(text="🔙 К статусу", callback_data="menu_status"))
|
|
||||||
try:
|
|
||||||
await message.edit_text("\n".join(lines), reply_markup=builder.as_markup())
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Proxy check result edit failed: {e}")
|
|
||||||
|
|
||||||
# =========================================================================
|
# =========================================================================
|
||||||
# УПРАВЛЕНИЕ ЗАДАЧЕЙ
|
# УПРАВЛЕНИЕ ЗАДАЧЕЙ
|
||||||
# =========================================================================
|
# =========================================================================
|
||||||
@@ -1873,7 +1865,8 @@ class BotInterface:
|
|||||||
InlineKeyboardButton(text="🗑️ Удалить", callback_data=f"tdelete_{task_id}")
|
InlineKeyboardButton(text="🗑️ Удалить", callback_data=f"tdelete_{task_id}")
|
||||||
)
|
)
|
||||||
|
|
||||||
builder.row(InlineKeyboardButton(text="🔙 К списку", callback_data="menu_streamers"))
|
back_target = "menu_tasks" if params.task_type in ("visit", "user_visit") else "menu_streamers"
|
||||||
|
builder.row(InlineKeyboardButton(text="🔙 К списку", callback_data=back_target))
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await callback.message.edit_text(text, reply_markup=builder.as_markup())
|
await callback.message.edit_text(text, reply_markup=builder.as_markup())
|
||||||
@@ -2184,15 +2177,15 @@ class BotInterface:
|
|||||||
await self._track_msg(params.user_id, msg.message_id)
|
await self._track_msg(params.user_id, msg.message_id)
|
||||||
|
|
||||||
async def _run_twitch(self, task_id: str, params: TaskParams, message: Message):
|
async def _run_twitch(self, task_id: str, params: TaskParams, message: Message):
|
||||||
from services.irc_service import TwitchIRCClient
|
from services.gql_chat_service import TwitchGQLChatClient
|
||||||
|
|
||||||
irc = TwitchIRCClient(params.channel, params.target_username)
|
irc = TwitchGQLChatClient(params.channel, params.target_username)
|
||||||
url_tasks: set = set()
|
url_tasks: set = set()
|
||||||
|
|
||||||
async def on_url(url, username):
|
async def on_url(url, username):
|
||||||
if params.stream_offline:
|
if params.stream_offline:
|
||||||
return
|
return
|
||||||
task = asyncio.create_task(self._process_url(url, username, params, message))
|
task = asyncio.create_task(self._process_url(url, username, task_id, params, message))
|
||||||
url_tasks.add(task)
|
url_tasks.add(task)
|
||||||
task.add_done_callback(url_tasks.discard)
|
task.add_done_callback(url_tasks.discard)
|
||||||
task.add_done_callback(
|
task.add_done_callback(
|
||||||
@@ -2206,10 +2199,10 @@ class BotInterface:
|
|||||||
await self.task_manager.pause_task(task_id)
|
await self.task_manager.pause_task(task_id)
|
||||||
await send_message_safe(
|
await send_message_safe(
|
||||||
message.bot, params.chat_id,
|
message.bot, params.chat_id,
|
||||||
f"⚠️ IRC #{params.channel}: ошибка соединения (3 попытки) — задача поставлена на паузу"
|
f"⚠️ GQL #{params.channel}: ошибка соединения (3 попытки) — задача поставлена на паузу"
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info(f"🚀 IRC 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))
|
||||||
|
|
||||||
@@ -2253,7 +2246,7 @@ class BotInterface:
|
|||||||
# Resume < 5 мин: продолжает с N+1 без задержки
|
# Resume < 5 мин: продолжает с N+1 без задержки
|
||||||
# Resume > 5 мин: серия этой ссылки отбрасывается, ждём следующий цикл
|
# Resume > 5 мин: серия этой ссылки отбрасывается, ждём следующий цикл
|
||||||
|
|
||||||
remaining_clicks = 0
|
remaining_clicks = params.remaining_clicks # восстанавливаем из сохранённого состояния
|
||||||
series_paused_at: Optional[float] = None
|
series_paused_at: Optional[float] = None
|
||||||
SERIES_EXPIRY = 5 * 60
|
SERIES_EXPIRY = 5 * 60
|
||||||
|
|
||||||
@@ -2301,6 +2294,7 @@ class BotInterface:
|
|||||||
break
|
break
|
||||||
|
|
||||||
# --- Серия кликов ---
|
# --- Серия кликов ---
|
||||||
|
loop = asyncio.get_event_loop()
|
||||||
for click_num in range(series_size):
|
for click_num in range(series_size):
|
||||||
if params.stopped or params.successful_visits >= max_v:
|
if params.stopped or params.successful_visits >= max_v:
|
||||||
break
|
break
|
||||||
@@ -2311,6 +2305,11 @@ class BotInterface:
|
|||||||
await send_message_safe(bot, params.chat_id, "⚠️ Баланс исчерпан — задача остановлена")
|
await send_message_safe(bot, params.chat_id, "⚠️ Баланс исчерпан — задача остановлена")
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# Фиксируем старт тика до посещения
|
||||||
|
is_last = (click_num == series_size - 1)
|
||||||
|
click_delay = 0 if is_last else params.get_click_delay()
|
||||||
|
tick_start = loop.time()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
visitor = self.browser_pool or self.browser_service
|
visitor = self.browser_pool or self.browser_service
|
||||||
result = await visitor.visit_page(params.url, params.get_reading_time())
|
result = await visitor.visit_page(params.url, params.get_reading_time())
|
||||||
@@ -2319,6 +2318,8 @@ class BotInterface:
|
|||||||
params.successful_visits += 1
|
params.successful_visits += 1
|
||||||
await self.balance_storage.deduct(params.user_id, 1)
|
await self.balance_storage.deduct(params.user_id, 1)
|
||||||
logger.info(f"✅ Visit {params.successful_visits}/{max_v}: {params.url[:40]}")
|
logger.info(f"✅ Visit {params.successful_visits}/{max_v}: {params.url[:40]}")
|
||||||
|
if self.storage:
|
||||||
|
await self.storage.save_task(task_id, params)
|
||||||
else:
|
else:
|
||||||
logger.warning(f"❌ Visit failed ({params.total_visits}): {result.error or 'unknown'}")
|
logger.warning(f"❌ Visit failed ({params.total_visits}): {result.error or 'unknown'}")
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
@@ -2330,14 +2331,18 @@ class BotInterface:
|
|||||||
# Клик завершён — проверяем паузу
|
# Клик завершён — проверяем паузу
|
||||||
if params.paused or params.stopped:
|
if params.paused or params.stopped:
|
||||||
remaining_clicks = series_size - click_num - 1
|
remaining_clicks = series_size - click_num - 1
|
||||||
|
params.remaining_clicks = remaining_clicks
|
||||||
if remaining_clicks > 0 and not params.stopped:
|
if remaining_clicks > 0 and not params.stopped:
|
||||||
series_paused_at = datetime.now().timestamp()
|
series_paused_at = datetime.now().timestamp()
|
||||||
logger.info(f"⏸️ Paused after click {click_num+1}/{series_size}, {remaining_clicks} saved")
|
logger.info(f"⏸️ Paused after click {click_num+1}/{series_size}, {remaining_clicks} saved")
|
||||||
break
|
break
|
||||||
|
|
||||||
# Задержка между кликами (кроме последнего)
|
# Задержка между кликами строго по таймеру: спим только остаток
|
||||||
if click_num < series_size - 1 and params.successful_visits < max_v:
|
if not is_last and params.successful_visits < max_v:
|
||||||
await _sleep_interruptible(params.get_click_delay())
|
elapsed = loop.time() - tick_start
|
||||||
|
remaining_wait = int(click_delay - elapsed)
|
||||||
|
if remaining_wait > 0:
|
||||||
|
await _sleep_interruptible(remaining_wait)
|
||||||
|
|
||||||
await self.task_manager.complete_task(task_id)
|
await self.task_manager.complete_task(task_id)
|
||||||
remaining = await self.balance_storage.get_balance(params.user_id)
|
remaining = await self.balance_storage.get_balance(params.user_id)
|
||||||
@@ -2354,7 +2359,7 @@ class BotInterface:
|
|||||||
if self.storage:
|
if self.storage:
|
||||||
await self.storage.save_task(task_id, params)
|
await self.storage.save_task(task_id, params)
|
||||||
|
|
||||||
async def _process_url(self, url: str, username: str, params: TaskParams, message: Message):
|
async def _process_url(self, url: str, username: str, task_id: str, params: TaskParams, message: Message):
|
||||||
"""Обрабатывает найденную ссылку (выполняется параллельно)."""
|
"""Обрабатывает найденную ссылку (выполняется параллельно)."""
|
||||||
try:
|
try:
|
||||||
params.links_found += 1
|
params.links_found += 1
|
||||||
@@ -2427,6 +2432,15 @@ class BotInterface:
|
|||||||
result = await visitor.visit_page(url, reading)
|
result = await visitor.visit_page(url, reading)
|
||||||
params.total_visits += 1
|
params.total_visits += 1
|
||||||
params.pending_visits = max(0, params.pending_visits - 1)
|
params.pending_visits = max(0, params.pending_visits - 1)
|
||||||
|
if result.error == "PROXY_UNAVAILABLE":
|
||||||
|
params.paused = True
|
||||||
|
params.auto_paused = True
|
||||||
|
await self.task_manager.update_task(task_id, paused=True)
|
||||||
|
await send_message_safe(
|
||||||
|
message.bot, params.chat_id,
|
||||||
|
"⚠️ Прокси недоступен — задача поставлена на паузу"
|
||||||
|
)
|
||||||
|
return
|
||||||
if result.success:
|
if result.success:
|
||||||
params.successful_visits += 1
|
params.successful_visits += 1
|
||||||
remaining -= 1
|
remaining -= 1
|
||||||
@@ -2952,18 +2966,21 @@ class BotInterface:
|
|||||||
|
|
||||||
cost = int(value * settings.CLICK_PRICE_RUB)
|
cost = int(value * settings.CLICK_PRICE_RUB)
|
||||||
|
|
||||||
# Сохраняем для шага подтверждения
|
|
||||||
self._payment_confirm[user_id] = {"clicks": value, "rub": cost}
|
self._payment_confirm[user_id] = {"clicks": value, "rub": cost}
|
||||||
|
rub_bal = await self.rub_storage.get_balance(user_id)
|
||||||
|
|
||||||
builder = InlineKeyboardBuilder()
|
builder = InlineKeyboardBuilder()
|
||||||
builder.row(
|
if rub_bal >= cost:
|
||||||
InlineKeyboardButton(text="✅ Подтвердить", callback_data="confirm_buy"),
|
builder.row(InlineKeyboardButton(
|
||||||
InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_confirm_buy"),
|
text=f"💰 С баланса ({rub_bal} ₽)",
|
||||||
)
|
callback_data="pay_from_balance"
|
||||||
|
))
|
||||||
|
builder.row(InlineKeyboardButton(text="💳 Оплатить криптой", callback_data="confirm_buy"))
|
||||||
|
builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_confirm_buy"))
|
||||||
await _edit_prompt(
|
await _edit_prompt(
|
||||||
f"🛒 Подтвердите покупку\n\n"
|
f"🛒 Подтвердите покупку\n\n"
|
||||||
f"🔢 {value} переходов\n"
|
f"🔢 {value} переходов × {fmt_price(settings.CLICK_PRICE_RUB)} ₽ = {cost} ₽"
|
||||||
f"💰 Стоимость: {cost} ₽",
|
+ (f"\n\n🪙 Ваш баланс: {rub_bal} ₽" if rub_bal > 0 else ""),
|
||||||
builder.as_markup()
|
builder.as_markup()
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -3102,6 +3119,7 @@ class BotInterface:
|
|||||||
text=f"📋 Задачи ({len(user_tasks)})",
|
text=f"📋 Задачи ({len(user_tasks)})",
|
||||||
callback_data=f"utasks_{target_user_id}",
|
callback_data=f"utasks_{target_user_id}",
|
||||||
))
|
))
|
||||||
|
builder.row(InlineKeyboardButton(text="✉️ Написать сообщение", callback_data=f"umsg_{target_user_id}"))
|
||||||
builder.row(InlineKeyboardButton(text="🔙 К списку", callback_data="menu_users"))
|
builder.row(InlineKeyboardButton(text="🔙 К списку", callback_data="menu_users"))
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
+17
-7
@@ -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:
|
||||||
|
|||||||
@@ -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 (
|
||||||
|
|||||||
+85
-85
@@ -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(
|
||||||
@@ -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)
|
||||||
|
|||||||
@@ -0,0 +1,264 @@
|
|||||||
|
"""
|
||||||
|
Twitch чат через IRC over WebSocket (wss://irc-ws.chat.twitch.tv:443).
|
||||||
|
Замена raw TCP IRC — тот же протокол, WebSocket транспорт.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
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()
|
||||||
|
self.irc_username = irc_username or settings.IRC_USERNAME
|
||||||
|
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
|
||||||
|
|
||||||
|
return stats
|
||||||
|
|
||||||
|
async def disconnect(self) -> None:
|
||||||
|
await self._cleanup()
|
||||||
@@ -173,7 +173,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 +181,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
@@ -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,
|
||||||
|
|||||||
@@ -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()
|
||||||
Reference in New Issue
Block a user