Compare commits

..

10 Commits

Author SHA1 Message Date
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
11 changed files with 984 additions and 290 deletions
+31
View File
@@ -1,8 +1,10 @@
import asyncio
import hashlib
import json
import secrets
import logging
from datetime import datetime, timedelta
from pathlib import Path
from typing import Dict, Optional
from config.settings import settings
@@ -10,6 +12,7 @@ from config.settings import settings
logger = logging.getLogger(__name__)
ADMIN_PASSWORD = settings.ADMIN_PASSWORD
_SESSIONS_FILE = Path("data/sessions.json")
class AuthManager:
@@ -22,6 +25,7 @@ class AuthManager:
self._max_attempts = 5
self._lockout_duration = 300
self._session_timeout = timedelta(minutes=session_timeout_minutes)
self._restore_sessions()
@staticmethod
def _hash_password(password: str, salt: str = None) -> tuple:
@@ -37,6 +41,10 @@ class AuthManager:
def is_authenticated(self, user_id: int) -> bool:
if user_id not in self._sessions:
return False
# Обычные пользователи — сессия бесконечная
if self._roles.get(user_id) != "admin":
return True
# Админы — таймаут применяется
if datetime.now() - self._sessions[user_id] > self._session_timeout:
self.logout(user_id)
return False
@@ -49,14 +57,37 @@ class AuthManager:
self._sessions[user_id] = datetime.now()
self._roles[user_id] = "admin" if is_admin else "user"
self._failed_attempts.pop(user_id, None)
self._save_sessions()
logger.info(f"User {user_id} logged in as {self._roles[user_id]}")
def logout(self, user_id: int):
self._sessions.pop(user_id, None)
self._roles.pop(user_id, None)
self._failed_attempts.pop(user_id, None)
self._save_sessions()
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:
if user_id not in self._failed_attempts:
return False
+6 -5
View File
@@ -54,11 +54,12 @@ class Settings(BaseSettings):
SOCKET_TIMEOUT: int = 30
RECV_TIMEOUT: float = 5.0
# Proxy
PROXY_FILE: str = "proxies.txt"
PROXY_COOLDOWN_TIME: int = 30
PROXY_MAX_USAGE_BEFORE_COOLDOWN: int = 5
# Sticky proxy (PlainProxies / аналоги)
# При заполненном STICKY_PROXY_HOST rotating-прокси для браузера не используются
STICKY_PROXY_HOST: str = " " # res-unlimited-XXXX.plainproxies.com
STICKY_PROXY_PORT: int = 8080
STICKY_PROXY_USER: str = " " # логин из личного кабинета
STICKY_PROXY_PASS: str = " " # пароль из личного кабинета
# Screenshots
SCREENSHOTS_DIR: str = ""
+151 -133
View File
@@ -22,7 +22,6 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder, ReplyKeyboardBuilder
from config.settings import settings
from managers.background_tasks import BackgroundTaskManager
from managers.proxy_manager import ProxyManager
from managers.task_manager import TaskManager, TaskParams
from services.browser_service import BrowserService
from services.browser_pool import BrowserPool
@@ -45,7 +44,6 @@ class BotInterface:
def __init__(
self,
background_tasks: BackgroundTaskManager,
proxy_manager: ProxyManager,
browser_service: BrowserService,
browser_pool: BrowserPool = None,
storage: TaskStorage = None,
@@ -53,7 +51,6 @@ class BotInterface:
bot_ref=None,
):
self.background_tasks = background_tasks
self.proxy_manager = proxy_manager
self.browser_service = browser_service
self.browser_pool = browser_pool
self.task_manager = TaskManager()
@@ -72,6 +69,7 @@ class BotInterface:
self._topup_confirm: Dict[int, int] = {} # user_id -> rub amount pending confirm
self._user_task_state: Dict[int, dict] = {} # шаги создания задачи пользователем
self._admin_balance_state: Dict[int, int] = {} # admin_id -> target_user_id
self._admin_msg_state: Dict[int, int] = {} # admin_id -> target_user_id
self._menu_msg: Dict[int, int] = {} # user_id -> inline keyboard message_id
self._menu_top_msg: Dict[int, int] = {} # user_id -> reply keyboard message_id
self._chat_history: Dict[int, list] = {}
@@ -258,29 +256,6 @@ class BotInterface:
await interface._show_status(callback.message, edit=True)
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")
async def cb_restart_confirm(callback: CallbackQuery):
if not await require_admin(callback):
@@ -508,8 +483,13 @@ class BotInterface:
if not await require_admin(callback):
return
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 callback.answer("🗑️ Удалена")
if is_visit:
await interface._show_tasks_list(callback.message, edit=True)
else:
await interface._show_streamers_list(callback.message, edit=True)
@@ -1239,6 +1219,28 @@ class BotInterface:
await interface._show_user_detail(callback, target_uid)
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_"))
async def cb_user_balance(callback: CallbackQuery):
if not await require_admin(callback):
@@ -1274,6 +1276,25 @@ class BotInterface:
user_id = message.from_user.id
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 not text:
@@ -1674,9 +1695,67 @@ class BotInterface:
await self._edit_or_send(message, text, builder.as_markup(), edit)
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()
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_main"))
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):
stats = await self.task_manager.get_stats()
proxy_stats = self.proxy_manager.get_stats()
bg_count = self.background_tasks.active_count
delta = datetime.now() - self._start_time
@@ -1718,14 +1796,11 @@ class BotInterface:
f"├ 🌐 Визиты: {stats['visit']}\n"
f"├ ⏸️ Пауза: {stats['paused']}\n"
f"└ 🔄 Активно: {stats['active']}\n\n"
f"🔌 Прокси: {proxy_stats['total']} / {proxy_stats['available']} доступно\n"
f"🔁 Фоновых задач: {bg_count}\n\n"
f"{sys_info}"
)
builder = InlineKeyboardBuilder()
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="menu_main"))
await self._edit_or_send(message, text, builder.as_markup(), edit)
@@ -1743,89 +1818,6 @@ class BotInterface:
await asyncio.sleep(0.5)
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}")
)
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:
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)
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()
async def on_url(url, username):
if params.stream_offline:
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)
task.add_done_callback(url_tasks.discard)
task.add_done_callback(
@@ -2206,10 +2199,10 @@ class BotInterface:
await self.task_manager.pause_task(task_id)
await send_message_safe(
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))
@@ -2253,7 +2246,7 @@ class BotInterface:
# Resume < 5 мин: продолжает с N+1 без задержки
# Resume > 5 мин: серия этой ссылки отбрасывается, ждём следующий цикл
remaining_clicks = 0
remaining_clicks = params.remaining_clicks # восстанавливаем из сохранённого состояния
series_paused_at: Optional[float] = None
SERIES_EXPIRY = 5 * 60
@@ -2301,6 +2294,7 @@ class BotInterface:
break
# --- Серия кликов ---
loop = asyncio.get_event_loop()
for click_num in range(series_size):
if params.stopped or params.successful_visits >= max_v:
break
@@ -2311,6 +2305,11 @@ class BotInterface:
await send_message_safe(bot, params.chat_id, "⚠️ Баланс исчерпан — задача остановлена")
return
# Фиксируем старт тика до посещения
is_last = (click_num == series_size - 1)
click_delay = 0 if is_last else params.get_click_delay()
tick_start = loop.time()
try:
visitor = self.browser_pool or self.browser_service
result = await visitor.visit_page(params.url, params.get_reading_time())
@@ -2319,6 +2318,8 @@ class BotInterface:
params.successful_visits += 1
await self.balance_storage.deduct(params.user_id, 1)
logger.info(f"✅ Visit {params.successful_visits}/{max_v}: {params.url[:40]}")
if self.storage:
await self.storage.save_task(task_id, params)
else:
logger.warning(f"❌ Visit failed ({params.total_visits}): {result.error or 'unknown'}")
except asyncio.CancelledError:
@@ -2330,14 +2331,18 @@ class BotInterface:
# Клик завершён — проверяем паузу
if params.paused or params.stopped:
remaining_clicks = series_size - click_num - 1
params.remaining_clicks = remaining_clicks
if remaining_clicks > 0 and not params.stopped:
series_paused_at = datetime.now().timestamp()
logger.info(f"⏸️ Paused after click {click_num+1}/{series_size}, {remaining_clicks} saved")
break
# Задержка между кликами (кроме последнего)
if click_num < series_size - 1 and params.successful_visits < max_v:
await _sleep_interruptible(params.get_click_delay())
# Задержка между кликами строго по таймеру: спим только остаток
if not is_last and params.successful_visits < max_v:
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)
remaining = await self.balance_storage.get_balance(params.user_id)
@@ -2354,7 +2359,7 @@ class BotInterface:
if self.storage:
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:
params.links_found += 1
@@ -2427,6 +2432,15 @@ class BotInterface:
result = await visitor.visit_page(url, reading)
params.total_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:
params.successful_visits += 1
remaining -= 1
@@ -2952,18 +2966,21 @@ class BotInterface:
cost = int(value * settings.CLICK_PRICE_RUB)
# Сохраняем для шага подтверждения
self._payment_confirm[user_id] = {"clicks": value, "rub": cost}
rub_bal = await self.rub_storage.get_balance(user_id)
builder = InlineKeyboardBuilder()
builder.row(
InlineKeyboardButton(text="✅ Подтвердить", callback_data="confirm_buy"),
InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_confirm_buy"),
)
if rub_bal >= cost:
builder.row(InlineKeyboardButton(
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(
f"🛒 Подтвердите покупку\n\n"
f"🔢 {value} переходов\n"
f"💰 Стоимость: {cost}",
f"🔢 {value} переходов × {fmt_price(settings.CLICK_PRICE_RUB)} ₽ = {cost}"
+ (f"\n\n🪙 Ваш баланс: {rub_bal}" if rub_bal > 0 else ""),
builder.as_markup()
)
@@ -3102,6 +3119,7 @@ class BotInterface:
text=f"📋 Задачи ({len(user_tasks)})",
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"))
try:
+34 -23
View File
@@ -10,7 +10,6 @@ from aiogram.types import BotCommand
from aiogram.client.session.aiohttp import AiohttpSession
from config.settings import settings
from core.logger import setup_logger
from managers.proxy_manager import ProxyManager
from managers.background_tasks import BackgroundTaskManager
from services.browser_service import BrowserService
from services.browser_pool import BrowserPool
@@ -43,10 +42,10 @@ class BotApplication:
def __init__(self):
self.bot: Bot = None
self.dispatcher: Dispatcher = None
self._initialized = False # True только после успешного _restore_tasks
self.proxy_manager = ProxyManager()
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.storage = TaskStorage()
self.chat_storage = ChatStorage()
@@ -55,7 +54,6 @@ class BotApplication:
self.interface = BotInterface(
background_tasks=self.background_tasks,
proxy_manager=self.proxy_manager,
browser_service=self.browser_service,
browser_pool=self.browser_pool,
storage=self.storage,
@@ -85,9 +83,6 @@ class BotApplication:
self.interface.register(self.dispatcher)
await self._update_bot_commands()
# Определяем типы прокси без явного протокола
await self.proxy_manager.detect_types()
# Запускаем пул браузеров
await self.browser_pool.start()
logger.info(f"Browser pool: {settings.BROWSER_POOL_SIZE} workers")
@@ -107,9 +102,24 @@ class BotApplication:
# Восстанавливаем задачи
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!")
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):
"""Восстанавливает задачи после перезапуска."""
saved_tasks = await self.storage.load_tasks()
@@ -119,19 +129,19 @@ class BotApplication:
return
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():
# Завершённые и остановленные — загружаем в память но не запускаем
# Завершённые и остановленные — только в память, не запускаем
if params.completed or params.stopped:
await self.interface.task_manager.add_task(task_id, params)
continue
# Активные задачи ставим на паузу и перезапускаем
# Если стрим был оффлайн — сохраняем этот статус
params.paused = True
was_paused = params.paused # сохраняем намерение пользователя
if params.stream_offline:
logger.info(f"Task {task_id}: stream was offline, keeping paused")
await self.interface.task_manager.add_task(task_id, params)
if params.task_type == "twitch_irc" and params.chat_id:
@@ -156,17 +166,16 @@ class BotApplication:
restored += 1
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:
await self.bot.send_message(
chat_id,
"🔄 Бот был перезапущен\n\n"
"Ваши задачи приостановлены.\n"
"Нажмите ▶️ в списке задач для возобновления."
)
if was_paused:
msg = "🔄 Бот был перезапущен\n\n⏸ Ваша задача остаётся на паузе."
else:
msg = "🔄 Бот был перезапущен\n\n▶️ Ваши задачи продолжают работу автоматически."
await self.bot.send_message(chat_id, msg)
except Exception as e:
logger.warning(f"Failed to notify user {uid}: {e}")
@@ -182,7 +191,9 @@ class BotApplication:
"""Остановка с сохранением."""
logger.info("Shutting down...")
# Сохраняем задачи
# Сохраняем только если инициализация прошла успешно —
# иначе task_manager пустой и мы затрём уже сохранённые задачи
if self._initialized:
tasks = self.interface.task_manager._tasks
await self.storage.save_tasks(tasks)
+17 -7
View File
@@ -54,6 +54,7 @@ class TaskStorage:
"max_ctr": params.max_ctr,
"started_at": params.started_at.isoformat() if params.started_at else None,
"stream_offline": params.stream_offline,
"remaining_clicks": params.remaining_clicks,
# pending_visits и total_planned_visits — runtime, не сохраняем
}
@@ -88,6 +89,7 @@ class TaskStorage:
min_ctr=data.get("min_ctr", 0.8),
max_ctr=data.get("max_ctr", 1.0),
stream_offline=data.get("stream_offline", False),
remaining_clicks=data.get("remaining_clicks", 0),
)
if data.get("started_at"):
try:
@@ -133,15 +135,23 @@ class TaskStorage:
return {}
async def save_task(self, task_id: str, params: TaskParams) -> None:
tasks = await self.load_tasks()
tasks[task_id] = params
await self.save_tasks(tasks)
async with self._lock:
try:
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:
tasks = await self.load_tasks()
if task_id in tasks:
del tasks[task_id]
await self.save_tasks(tasks)
async with self._lock:
try:
data = await asyncio.to_thread(self._load_sync)
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:
+18 -7
View File
@@ -59,6 +59,7 @@ class TaskParams:
links_found: int = 0
pending_visits: int = 0 # кликов в очереди (runtime, уменьшается)
total_planned_visits: int = 0 # всего запланировано кликов по всем ссылкам (runtime)
remaining_clicks: int = 0 # остаток кликов в текущей серии (сохраняется)
chat_id: Optional[int] = None
user_id: Optional[int] = None # Telegram user_id, если задача создана пользователем
@@ -271,8 +272,12 @@ class TaskManager:
return await self.get_tasks_by_type("twitch_irc")
async def get_visit_tasks(self) -> Dict[str, TaskParams]:
"""Только задачи посещения."""
return await self.get_tasks_by_type("visit")
"""Задачи посещения (admin и user)."""
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:
"""Общая статистика."""
@@ -280,7 +285,7 @@ class TaskManager:
for p in self._tasks.values():
if p.task_type == "twitch_irc":
twitch += 1
elif p.task_type == "visit":
elif p.task_type in ("visit", "user_visit"):
visit += 1
if p.paused:
paused += 1
@@ -299,12 +304,18 @@ class TaskManager:
status = params.get_status_emoji()
if params.task_type == "twitch_irc":
user_tag = f" | 👤uid:{params.user_id}" if params.user_id else ""
return (
f"{status} 📺 **{params.channel}**\n"
f" `{task_id[:12]}...`\n"
f" 👤 @{params.target_username} | "
f"🔗{params.links_found} | "
f"📊{params.total_visits}\n"
f" `{task_id[:12]}...`{user_tag}\n"
f" 🔗{params.links_found} | 📊{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:
return (
+76 -76
View File
@@ -1,5 +1,5 @@
"""
Browser service с поддержкой SOCKS5 через локальный HTTP туннель.
Browser service с поддержкой HTTP и sticky-прокси.
"""
import asyncio
@@ -12,11 +12,9 @@ from dataclasses import dataclass
from camoufox import DefaultAddons
from camoufox.async_api import AsyncCamoufox
from camoufox.exceptions import InvalidIP
from camoufox.exceptions import InvalidIP, InvalidProxy
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__)
@@ -73,6 +71,44 @@ def _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:
"""Возвращает BCP47 locale без script-тега по реальному IP прокси."""
try:
@@ -104,14 +140,15 @@ class VisitResult:
class BrowserService:
"""
Сервис для посещения страниц через Camoufox.
Поддерживает HTTP, SOCKS5 (через туннель) и прямое соединение.
Поддерживает HTTP и sticky-прокси.
"""
def __init__(self, proxy_manager: ProxyManager):
self.proxy_manager = proxy_manager
self.socks5_pool = Socks5ProxyPool(idle_timeout=300)
_PROXY_FAIL_THRESHOLD = 5 # пауза задачи после N подряд неудачных визитов
def __init__(self):
self.screenshots_dir = Path(settings.SCREENSHOTS_DIR)
self.screenshots_dir.mkdir(exist_ok=True)
self._proxy_fail_streak: int = 0
async def visit_page(
self,
@@ -125,42 +162,22 @@ class BrowserService:
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):
proxy = None
use_socks5 = False
try:
if self.proxy_manager.has_proxies():
proxy = await self.proxy_manager.acquire_proxy(timeout=30)
proxy_locale = None
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]
proxy_config, proxy_url, proxy_info = _make_sticky_proxy_config()
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)
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:
logger.info(f"Using sticky proxy: {proxy_info} ip={real_ip}")
camoufox_kwargs = dict(
headless=True,
geoip=True,
@@ -170,56 +187,45 @@ class BrowserService:
)
if proxy_locale:
camoufox_kwargs["locale"] = proxy_locale
async with AsyncCamoufox(**camoufox_kwargs) as browser:
result = await self._browse_page(
browser, url,
proxy.id if proxy else "direct",
reading_time
)
else:
return VisitResult(url=url, success=False, error="No proxy available")
# Определяем тип ошибки: ошибка соединения (прокси сломан) vs ошибка загрузки (прокси ok)
async def _run_browser():
async with AsyncCamoufox(**camoufox_kwargs) as browser:
return await self._browse_page(browser, url, proxy_info, reading_time)
result = await asyncio.wait_for(_run_browser(), timeout=attempt_timeout)
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(
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:
logger.warning(f"Visit attempt {attempt + 1} proxy connection failed, retrying with new proxy")
continue
elif is_load_error:
# Ошибка загрузки — не меняем прокси, просто возвращаем результат
# (сайт мог быть временно недоступен, прокси ок)
logger.warning(f"Visit attempt {attempt + 1} page load error: {result.error or 'unknown'}")
self._proxy_fail_streak = 0
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}")
if proxy:
await self.proxy_manager.release_proxy(proxy, success=False)
if use_socks5:
await self.socks5_pool.release(proxy)
except Exception as e:
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))
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)")
async def _browse_page(
@@ -373,13 +379,7 @@ class BrowserService:
return None
async def cleanup(self):
"""Очистка с таймаутом."""
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}")
pass
async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult:
return await self.visit_page(url, reading_time)
+264
View File
@@ -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()
+7 -4
View File
@@ -173,7 +173,6 @@ class TwitchIRCClient:
) -> dict:
stats = {"links_found": 0, "messages": 0, "reconnects": 0}
start_time = datetime.now()
last_message_time = datetime.now()
drain_errors = 0
@@ -182,20 +181,24 @@ class TwitchIRCClient:
logger.info(f"⏸️ #{self.channel} waiting for resume before connecting...")
while is_active and not is_active():
await asyncio.sleep(2)
if not is_active and not is_active():
return stats
logger.info(f"▶️ #{self.channel} resumed, connecting to IRC")
if not await self.connect():
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():
logger.info(f"⏸️ #{self.channel} paused, waiting for resume...")
await self.disconnect()
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} resumed, reconnecting")
await self.connect()
last_message_time = datetime.now()
+13 -4
View File
@@ -56,6 +56,7 @@ class VisitScheduler:
Statistics dict
"""
total_reading_time = 0
loop = asyncio.get_event_loop()
try:
while True:
@@ -63,15 +64,18 @@ class VisitScheduler:
logger.info(f"Visit limit reached: {self.max_visits}")
break
# Randomize parameters
# Randomize parameters once per cycle
current_reading = random.randint(self.min_reading, self.max_reading)
current_delay = random.randint(self.min_delay, self.max_delay)
self.visit_count += 1
# Record start time so the next visit fires strictly on the timer
tick_start = loop.time()
logger.info(
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
@@ -94,11 +98,16 @@ class VisitScheduler:
if self.max_visits and self.visit_count >= self.max_visits:
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:
await self.on_progress(self.visit_count, current_delay)
await asyncio.sleep(current_delay)
if remaining > 0:
await asyncio.sleep(remaining)
except asyncio.CancelledError:
logger.info("Visit scheduler cancelled")
+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()