init auth and main function

This commit is contained in:
Yuriy Yuriev
2026-05-14 16:15:58 +07:00
commit 27c53e897e
47 changed files with 5750 additions and 0 deletions
+2
View File
@@ -0,0 +1,2 @@
from services.browser_service import BrowserService, VisitResult
from services.socks5_to_http_proxy import Socks5ToHttpProxy, Socks5ProxyPool
+450
View File
@@ -0,0 +1,450 @@
"""
Browser service с поддержкой SOCKS5 через локальный HTTP туннель.
"""
import asyncio
import random
import logging
from pathlib import Path
from datetime import datetime
from typing import Optional
from dataclasses import dataclass
from camoufox import DefaultAddons
from camoufox.async_api import AsyncCamoufox
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__)
@dataclass
class VisitResult:
"""Результат посещения страницы."""
url: str
initial_url: Optional[str] = None
final_url: Optional[str] = None
redirect_occurred: bool = False
proxy_info: str = "direct"
reading_time: int = 0
screenshot_path: Optional[str] = None
success: bool = False
error: Optional[str] = None
class BrowserService:
"""
Сервис для посещения страниц через Camoufox.
Поддерживает HTTP, SOCKS5 (через туннель) и прямое соединение.
"""
def __init__(self, proxy_manager: ProxyManager):
self.proxy_manager = proxy_manager
self.socks5_pool = Socks5ProxyPool(idle_timeout=300)
self.screenshots_dir = Path(settings.SCREENSHOTS_DIR)
self.screenshots_dir.mkdir(exist_ok=True)
async def visit_page(
self,
url: str,
reading_time: Optional[int] = None
) -> VisitResult:
"""Посещает страницу."""
if reading_time is None:
reading_time = random.randint(
settings.DEFAULT_MIN_READING,
settings.DEFAULT_MAX_READING
)
proxy = None
use_socks5 = False
try:
# Получаем прокси
if self.proxy_manager.has_proxies():
proxy = await self.proxy_manager.acquire_proxy(timeout=30)
# Определяем конфигурацию
if proxy and proxy.proxy_type in [ProxyType.SOCKS5, ProxyType.SOCKS4]:
# SOCKS5 через локальный HTTP туннель
proxy_config = await self.socks5_pool.get_proxy_config(proxy)
use_socks5 = True
logger.info(f"Using SOCKS5 tunnel for: {proxy.id}")
elif proxy:
# HTTP прокси напрямую
proxy_config = proxy.proxy_config
logger.info(f"Using HTTP proxy: {proxy.id}")
else:
# Прямое соединение
proxy_config = None
logger.info("Direct connection")
# Запускаем браузер
if proxy_config:
async with AsyncCamoufox(
headless=False,
geoip=True,
humanize=True,
exclude_addons=[DefaultAddons.UBO],
proxy=proxy_config
) as browser:
result = await self._browse_page(
browser, url,
proxy.id if proxy else "direct",
reading_time
)
else:
async with AsyncCamoufox(
headless=False,
geoip=False,
humanize=True,
locale="ru-RU",
exclude_addons=[DefaultAddons.UBO]
) as browser:
result = await self._browse_page(
browser, url, "direct", reading_time
)
# Освобождаем
if proxy:
await self.proxy_manager.release_proxy(proxy, result.success)
if use_socks5:
await self.socks5_pool.release(proxy)
return result
except Exception as e:
logger.error(f"Visit error: {e}", exc_info=True)
if proxy:
await self.proxy_manager.release_proxy(proxy, success=False)
if use_socks5:
await self.socks5_pool.release(proxy)
return VisitResult(url=url, success=False, error=str(e))
async def _browse_page(
self, browser, url: str, proxy_info: str, reading_time: int
) -> VisitResult:
"""Выполняет просмотр страницы."""
page = await browser.new_page()
try:
await page.set_viewport_size({
"width": settings.VIEWPORT_WIDTH,
"height": settings.VIEWPORT_HEIGHT
})
goto_task = asyncio.create_task(
page.goto(url, wait_until="commit", timeout=30000)
)
await self._move_mouse(page, goto_task)
await goto_task
initial_url = page.url
final_url, redirect = await self._wait_redirect(page, initial_url)
await self._simulate_reading(page, reading_time)
screenshot_path = None
if settings.SCREENSHOTS_DIR:
screenshot_path = await self._take_screenshot(page, final_url)
return VisitResult(
url=url,
initial_url=initial_url,
final_url=final_url,
redirect_occurred=redirect,
proxy_info=proxy_info,
reading_time=reading_time,
screenshot_path=screenshot_path,
success=True
)
except Exception as e:
return VisitResult(url=url, reading_time=reading_time,
success=False, error=str(e))
finally:
if page.url == url:
logger.error()
await page.close()
async def _move_mouse(self, page, task):
"""Двигает мышь."""
w, h = settings.VIEWPORT_WIDTH, settings.VIEWPORT_HEIGHT
try:
while not task.done():
await page.mouse.move(
random.randint(100, w-100),
random.randint(100, h-100)
)
await asyncio.sleep(random.uniform(0.1, 0.3))
except:
pass
async def _wait_redirect(self, page, initial_url):
"""Ждет редирект."""
loop = asyncio.get_running_loop()
start = loop.time()
redirect = False
current = initial_url
while loop.time() - start < settings.MAX_REDIRECT_WAIT:
await asyncio.sleep(2)
new_url = page.url
if new_url != current:
current = new_url
redirect = True
task = asyncio.create_task(
page.wait_for_load_state("load", timeout=30000)
)
await self._move_mouse(page, task)
await task
break
await self._random_action(page)
return page.url, redirect
async def _simulate_reading(self, page, duration):
"""Симулирует чтение."""
loop = asyncio.get_running_loop()
start = loop.time()
while loop.time() - start < duration:
await self._random_action(page)
async def _random_action(self, page):
"""Случайное действие."""
try:
action = random.choice(['move', 'scroll', 'pause'])
if action == 'move':
x = random.randint(100, settings.VIEWPORT_WIDTH-100)
y = random.randint(100, settings.VIEWPORT_HEIGHT-100)
await page.mouse.move(x, y)
await asyncio.sleep(random.uniform(0.3, 1.0))
elif action == 'scroll':
await page.evaluate(f"window.scrollBy(0, {random.randint(100, 300)})")
await asyncio.sleep(random.uniform(0.3, 0.8))
else:
await asyncio.sleep(random.uniform(1, 3))
except:
await asyncio.sleep(0.5)
async def _take_screenshot(self, page, url):
"""Скриншот."""
try:
ts = datetime.now().strftime("%Y%m%d_%H%M%S")
domain = url.replace("https://", "").replace("http://", "").split("/")[0][:30]
fp = self.screenshots_dir / f"{ts}_{domain}.png"
await page.screenshot(path=str(fp), full_page=True)
return str(fp)
except Exception as e:
logger.error(f"Screenshot: {e}")
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}")
async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult:
if reading_time is None:
reading_time = random.randint(settings.DEFAULT_MIN_READING, settings.DEFAULT_MAX_READING)
proxy = None
for attempt in range(2): # Две попытки с разными прокси
try:
if self.proxy_manager.has_proxies():
proxy = await self.proxy_manager.acquire_proxy(timeout=90)
if proxy and proxy.proxy_type in [ProxyType.SOCKS5, ProxyType.SOCKS4]:
proxy_config = await self.socks5_pool.get_proxy_config(proxy)
elif proxy:
proxy_config = proxy.proxy_config
else:
proxy_config = None
browser_kwargs = {
"headless": False,
"geoip": True,
"humanize": True,
"exclude_addons": [DefaultAddons.UBO],
}
if proxy_config:
browser_kwargs["proxy"] = proxy_config
async with AsyncCamoufox(**browser_kwargs) as browser:
result = await self._browse_page(
browser,
url,
proxy.id if proxy else "direct",
reading_time
)
if proxy:
await self.proxy_manager.release_proxy(proxy, result.success)
if result.success:
return result
except Exception as e:
logger.warning(f"Attempt {attempt + 1} failed: {e}")
if proxy:
await self.proxy_manager.release_proxy(proxy, success=False)
await self.socks5_pool.release(proxy)
proxy = None
if attempt < 1:
await asyncio.sleep(2)
return VisitResult(url=url, success=False, error="All attempts failed")
async def _browse_twitch_ref(self, browser, url: str, channel: str, proxy_info: str, reading_time: int) -> VisitResult:
"""Посещает страницу как переход с Twitch канала."""
page = await browser.new_page()
twitch_url = f"https://www.twitch.tv/{channel}"
try:
await page.set_viewport_size({"width": settings.VIEWPORT_WIDTH, "height": settings.VIEWPORT_HEIGHT})
# Заголовки имитирующие переход с Twitch
await page.set_extra_http_headers({
"Referer": twitch_url,
"Origin": "https://www.twitch.tv",
})
logger.info(f"📺 From Twitch: {channel}{url}")
# Переход с Referer
goto = asyncio.create_task(page.goto(url, referer=twitch_url, wait_until="commit", timeout=30000))
await self._move_mouse(page, goto)
await goto
initial_url = page.url
final_url, redirect = await self._wait_redirect(page, initial_url)
await self._simulate_reading(page, reading_time)
return VisitResult(
url=url, initial_url=initial_url, final_url=final_url,
redirect_occurred=redirect, proxy_info=proxy_info,
reading_time=reading_time, success=True
)
except Exception as e:
return VisitResult(url=url, success=False, error=str(e))
finally:
await page.close()
async def _browse_from_twitch(
self,
browser,
url: str,
channel: str,
proxy_info: str,
reading_time: int
) -> VisitResult:
"""
Эмулирует переход с Twitch:
1. Открывает страницу стримера
2. Прокручивает чат
3. Переходит по ссылке
"""
page = await browser.new_page()
twitch_url = f"https://www.twitch.tv/{channel}"
try:
await page.set_viewport_size({
"width": settings.VIEWPORT_WIDTH,
"height": settings.VIEWPORT_HEIGHT
})
# === Шаг 1: Заходим на Twitch ===
logger.info(f"📺 Opening Twitch: {twitch_url}")
goto_task = asyncio.create_task(
page.goto(twitch_url, wait_until="commit", timeout=30000)
)
await self._move_mouse(page, goto_task)
await goto_task
# Ждем загрузки
await asyncio.sleep(random.uniform(2, 4))
# === Шаг 2: Имитируем просмотр стрима ===
logger.info("👀 Watching stream...")
# Прокручиваем страницу как будто смотрим
for _ in range(random.randint(2, 4)):
scroll = random.randint(200, 500)
await page.evaluate(f"window.scrollBy(0, {scroll})")
await asyncio.sleep(random.uniform(0.5, 1.5))
# Двигаем мышь (как будто читаем чат)
w, h = settings.VIEWPORT_WIDTH, settings.VIEWPORT_HEIGHT
for _ in range(random.randint(3, 6)):
x = random.randint(100, w - 100)
y = random.randint(100, h - 100)
await page.mouse.move(x, y)
await asyncio.sleep(random.uniform(0.3, 0.8))
# === Шаг 3: Переходим по ссылке ===
logger.info(f"🔗 Clicking link: {url}")
# Создаем новую вкладку для перехода (как target="_blank")
# Или просто переходим с Referer
await page.evaluate(f"""
window.open('{url}', '_blank');
""")
# Ждем открытия новой вкладки
await asyncio.sleep(2)
# Получаем новую вкладку
pages = await browser.pages()
if len(pages) > 1:
new_page = pages[-1]
await new_page.bring_to_front()
else:
# Если вкладка не открылась - переходим в этой же
new_page = page
await new_page.goto(url, referer=twitch_url, wait_until="commit")
initial_url = new_page.url
logger.info(f"📍 Landed: {initial_url}")
# === Шаг 4: Ждем редирект ===
final_url, redirect = await self._wait_redirect(new_page, initial_url)
# === Шаг 5: Читаем страницу ===
await self._simulate_reading(new_page, reading_time)
# Закрываем новую вкладку если она отдельная
if new_page != page:
await new_page.close()
# Возвращаемся на Twitch и закрываем
await page.close()
return VisitResult(
url=url,
initial_url=initial_url,
final_url=final_url,
redirect_occurred=redirect,
proxy_info=proxy_info,
reading_time=reading_time,
success=True
)
except Exception as e:
logger.error(f"Browse from Twitch error: {e}")
try:
await page.close()
except:
pass
return VisitResult(url=url, reading_time=reading_time, success=False, error=str(e))
+264
View File
@@ -0,0 +1,264 @@
"""
IRC сервис для мониторинга Twitch чата.
"""
import asyncio
import re
import logging
from typing import Optional, Callable, Awaitable, List
from datetime import datetime
from urllib.parse import urlparse
from config.settings import settings
logger = logging.getLogger(__name__)
class TwitchIRCClient:
"""
Twitch IRC клиент для мониторинга чата.
"""
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._reader: Optional[asyncio.StreamReader] = None
self._writer: Optional[asyncio.StreamWriter] = None
self._connected = False
self._url_pattern = re.compile(r'https?://[^\s<>"]+')
async def connect(self) -> bool:
"""Подключение к IRC серверу."""
try:
import ssl
ssl_context = ssl.create_default_context()
self._reader, self._writer = await asyncio.open_connection(
settings.IRC_SERVER, settings.IRC_PORT,
ssl=ssl_context
)
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"✅ Подключен к #{self.channel}")
return True
except Exception as e:
logger.error(f"❌ Ошибка подключения: {e}")
return False
async def _send(self, message: str) -> None:
"""Отправка сообщения."""
if self._writer:
self._writer.write(f"{message}\r\n".encode())
await self._writer.drain()
async def _read_line(self) -> Optional[str]:
"""Чтение строки с обработкой PING."""
if not self._reader:
return None
try:
line = await self._reader.readline()
if not line:
logger.warning("Empty line - connection closed")
self._connected = False
return None
decoded = line.decode('utf-8', errors='ignore').strip()
# PING - отвечаем СРАЗУ
if decoded.startswith('PING'):
pong_msg = f"PONG {decoded.split()[1]}"
self._writer.write(f"{pong_msg}\r\n".encode())
try:
await self._writer.drain()
except (ConnectionResetError, BrokenPipeError, OSError) as e:
logger.warning(f"Failed to send PONG: {e}")
logger.debug(f"🏓 {pong_msg}")
return f"PONG {decoded.split()[1]}" # Возвращаем PONG для отслеживания
return decoded
except Exception as e:
logger.error(f"Read error: {e}")
self._connected = False
return None
def _parse_message(self, line: str) -> Optional[dict]:
if not line or 'PRIVMSG' not in line:
return None
try:
# Парсим display-name из тегов
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', '')
# Извлекаем сообщение
if ' :' in line:
message_text = line.split(' :', 1)[1]
else:
message_text = ''
# Пропускаем пустые сообщения
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]:
"""Извлечение URL из текста."""
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, # Функция проверки активности
) -> dict:
stats = {"links_found": 0, "messages": 0, "reconnects": 0}
start_time = datetime.now()
last_message_time = datetime.now()
await self.connect()
while (datetime.now() - start_time).seconds < duration:
# === ПРОВЕРКА АКТИВНОСТИ ===
if is_active and not is_active():
logger.info("⏸️ Task paused/stopped, disconnecting...")
break
# === ПЕРЕПОДКЛЮЧЕНИЕ КАЖДЫЕ 30 СЕКУНД ===
if (datetime.now() - last_message_time).seconds > 30:
logger.info("🔄 Reconnecting (30s)...")
# Сначала читаем всё что осталось
try:
while True:
line = await asyncio.wait_for(self._reader.readline(), timeout=1.0)
if not line:
break
decoded = line.decode('utf-8', errors='ignore').strip()
if 'PRIVMSG' in decoded:
msg = self._parse_message(decoded)
if msg:
stats["messages"] += 1
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 Exception as e:
logger.warning(f"URL callback error: {e}")
except Exception as e:
logger.warning(f"Error draining IRC data: {e}")
# Переподключаемся
try:
await self.disconnect()
except Exception as e:
logger.warning(f"Disconnect error: {e}")
await asyncio.sleep(1)
await self.connect()
stats["reconnects"] += 1
last_message_time = datetime.now()
continue
# Читаем строку
try:
line = await asyncio.wait_for(self._read_line(), timeout=5.0)
except (asyncio.TimeoutError, asyncio.CancelledError):
continue
except Exception as e:
logger.warning(f"IRC read error: {e}")
continue
if not line:
continue
if line.startswith('PONG'):
continue
msg = self._parse_message(line)
if not msg:
continue
last_message_time = datetime.now() # Обновляем таймер
stats["messages"] += 1
logger.debug(f"IRC message 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:
pass
return stats
async def disconnect(self) -> None:
"""Закрытие соединения."""
self._connected = False
if self._writer:
try:
self._writer.close()
await self._writer.wait_closed()
except (ConnectionResetError, BrokenPipeError, OSError) as e:
logger.debug(f"Disconnect warning: {e}")
self._writer = None
self._reader = None
+381
View File
@@ -0,0 +1,381 @@
"""
SOCKS5 → HTTP/HTTPS прокси туннель с поддержкой CONNECT метода.
"""
import asyncio
import socket
import logging
from typing import Optional, Dict
import traceback
from aiohttp import ClientSession, ClientTimeout
from aiohttp_socks import ProxyConnector
logger = logging.getLogger(__name__)
class Socks5ToHttpProxy:
"""
Локальный HTTP/HTTPS → SOCKS5 прокси с поддержкой CONNECT.
"""
def __init__(
self,
socks5_host: str,
socks5_port: int,
username: Optional[str] = None,
password: Optional[str] = None
):
self.socks5_host = socks5_host
self.socks5_port = int(socks5_port)
self.username = username
self.password = password
self._session: Optional[ClientSession] = None
self._server = None
self._port: Optional[int] = None
self._active_connections = 0
self._lock = asyncio.Lock()
@property
def local_url(self) -> str:
return f"http://127.0.0.1:{self._port}"
@property
def proxy_config_for_browser(self) -> dict:
return {'server': self.local_url}
async def start(self) -> str:
"""Запускает прокси сервер."""
self._port = self._find_free_port()
socks_url = self._build_socks5_url()
connector = ProxyConnector.from_url(socks_url)
timeout = ClientTimeout(total=30, connect=15)
self._session = ClientSession(connector=connector, timeout=timeout)
self._server = await asyncio.start_server(
self._handle_client,
'127.0.0.1',
self._port
)
logger.info(f"🟢 SOCKS5: {self.local_url}")
return self.local_url
async def stop(self):
"""Останавливает прокси с таймаутом."""
logger.debug(f"Stopping proxy 127.0.0.1:{self._port}")
# Закрываем сервер (не принимать новые соединения)
if self._server:
self._server.close()
# Ждем завершения активных соединений с таймаутом
try:
async with asyncio.timeout(5): # Максимум 5 секунд на завершение
while self._active_connections > 0:
await asyncio.sleep(0.1)
except asyncio.TimeoutError:
logger.warning(f"Force stopping proxy with {self._active_connections} active connections")
# Закрываем сервер принудительно
if self._server:
try:
await asyncio.wait_for(self._server.wait_closed(), timeout=3)
except asyncio.TimeoutError:
logger.warning("Server close timeout")
self._server = None
# Закрываем сессию
if self._session:
try:
await asyncio.wait_for(self._session.close(), timeout=3)
except asyncio.TimeoutError:
pass
self._session = None
logger.debug(f"Proxy stopped: 127.0.0.1:{self._port}")
def _find_free_port(self) -> int:
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
s.bind(('127.0.0.1', 0))
return s.getsockname()[1]
def _build_socks5_url(self) -> str:
if self.username and self.password:
return f"socks5://{self.username}:{self.password}@{self.socks5_host}:{self.socks5_port}"
return f"socks5://{self.socks5_host}:{self.socks5_port}"
async def _handle_client(
self,
reader: asyncio.StreamReader,
writer: asyncio.StreamWriter
):
"""Обрабатывает входящее подключение."""
async with self._lock:
self._active_connections += 1
try:
# Таймаут на чтение первого запроса
request_line = await asyncio.wait_for(
reader.readline(),
timeout=10
)
if not request_line:
return
request_text = request_line.decode('utf-8', errors='ignore').strip()
if request_text.upper().startswith('CONNECT'):
await self._handle_connect(request_text, reader, writer)
else:
await self._handle_http(request_line, reader, writer)
except asyncio.TimeoutError:
logger.debug("Client timeout")
except ConnectionResetError:
logger.debug("Connection reset by client")
except Exception as e:
if "Server disconnected" not in str(e):
logger.error(f"Client error: {e}")
finally:
async with self._lock:
self._active_connections -= 1
try:
writer.close()
except Exception:
pass
async def _handle_connect(
self,
request_line: str,
reader: asyncio.StreamReader,
writer: asyncio.StreamWriter
):
"""Обрабатывает CONNECT запрос."""
import socks
remote_socket = None
try:
parts = request_line.split()
if len(parts) < 2:
writer.write(b'HTTP/1.1 400 Bad Request\r\n\r\n')
return
target = parts[1]
if ':' in target:
host, port = target.rsplit(':', 1)
port = int(port)
else:
host = target
port = 443
# Пропускаем заголовки (с таймаутом)
while True:
line = await asyncio.wait_for(reader.readline(), timeout=5)
if not line or line in [b'\r\n', b'\n']:
break
# Создаем SOCKS5 соединение
remote_socket = socks.socksocket()
remote_socket.settimeout(15)
if self.username:
remote_socket.set_proxy(
socks.SOCKS5, self.socks5_host, self.socks5_port,
username=self.username, password=self.password
)
else:
remote_socket.set_proxy(
socks.SOCKS5, self.socks5_host, self.socks5_port
)
# Подключаемся (в отдельном потоке чтобы не блокировать)
loop = asyncio.get_running_loop()
await loop.run_in_executor(None, remote_socket.connect, (host, port))
# Отвечаем клиенту
writer.write(b'HTTP/1.1 200 Connection Established\r\n\r\n')
await writer.drain()
# Пересылка данных с таймаутом
await asyncio.wait_for(
self._relay(reader, writer, remote_socket),
timeout=60
)
except asyncio.TimeoutError:
logger.debug("CONNECT timeout")
except Exception as e:
logger.debug(f"CONNECT error: {e}")
finally:
if remote_socket:
try:
remote_socket.close()
except Exception:
pass
async def _relay(
self,
client_reader: asyncio.StreamReader,
client_writer: asyncio.StreamWriter,
remote_socket
):
"""Двунаправленная пересылка данных."""
loop = asyncio.get_running_loop()
running = True
async def client_to_remote():
try:
while running:
data = await asyncio.wait_for(client_reader.read(8192), timeout=30)
if not data:
break
await loop.sock_sendall(remote_socket, data)
except Exception:
pass
async def remote_to_client():
try:
while running:
data = await loop.sock_recv(remote_socket, 8192)
if not data:
break
client_writer.write(data)
await client_writer.drain()
except Exception:
pass
try:
await asyncio.gather(client_to_remote(), remote_to_client())
except Exception:
pass
finally:
running = False
async def _handle_http(
self,
request_line: bytes,
reader: asyncio.StreamReader,
writer: asyncio.StreamWriter
):
"""Обрабатывает HTTP запрос."""
try:
first_line = request_line.decode('utf-8', errors='ignore').strip()
parts = first_line.split()
if len(parts) < 2:
return
method = parts[0]
path = parts[1]
# Читаем заголовки
headers = {}
while True:
line = await asyncio.wait_for(reader.readline(), timeout=5)
if not line or line in [b'\r\n', b'\n']:
break
decoded = line.decode('utf-8', errors='ignore').strip()
if ':' in decoded:
key, value = decoded.split(':', 1)
headers[key.strip()] = value.strip()
# Тело запроса
body = b''
content_length = int(headers.get('Content-Length', 0))
if content_length > 0:
body = await reader.readexactly(content_length)
host = headers.get('Host', 'localhost')
url = f"http://{host}{path}"
headers.pop('Proxy-Connection', None)
# Проксируем
async with self._session.request(
method=method, url=url,
headers=headers, data=body if body else None,
allow_redirects=True
) as response:
response_body = await response.read()
status_line = f"HTTP/1.1 {response.status} {response.reason}\r\n"
response_text = status_line
for key, value in response.headers.items():
if key.lower() != 'transfer-encoding':
response_text += f"{key}: {value}\r\n"
response_text += f"Content-Length: {len(response_body)}\r\n\r\n"
writer.write(response_text.encode() + response_body)
await writer.drain()
except asyncio.TimeoutError:
pass
except Exception as e:
if "Server disconnected" not in str(e):
logger.debug(f"HTTP error: {e}")
class Socks5ProxyPool:
"""Пул SOCKS5→HTTP прокси."""
def __init__(self, idle_timeout: int = 300):
self._proxies: Dict[str, Socks5ToHttpProxy] = {}
self._last_used: Dict[str, float] = {}
self._lock = asyncio.Lock()
self._idle_timeout = idle_timeout
self._cleanup_task: Optional[asyncio.Task] = None
async def get_proxy_config(self, socks5_proxy) -> dict:
"""Получает или создает конфигурацию HTTP прокси для SOCKS5."""
key = socks5_proxy.id
async with self._lock:
if key not in self._proxies:
proxy = Socks5ToHttpProxy(
socks5_host=socks5_proxy.ip,
socks5_port=int(socks5_proxy.port),
username=socks5_proxy.login,
password=socks5_proxy.password
)
await proxy.start()
self._proxies[key] = proxy
loop = asyncio.get_running_loop()
self._last_used[key] = loop.time()
return self._proxies[key].proxy_config_for_browser
async def release(self, socks5_proxy):
"""Помечает прокси как неиспользуемый."""
key = socks5_proxy.id
async with self._lock:
if key in self._last_used:
self._last_used[key] = 0
async def stop_all(self):
"""Останавливает все прокси с таймаутом."""
if self._cleanup_task:
self._cleanup_task.cancel()
async with self._lock:
proxies = list(self._proxies.values())
self._proxies.clear()
self._last_used.clear()
for proxy in proxies:
try:
await asyncio.wait_for(proxy.stop(), timeout=3)
except (asyncio.TimeoutError, Exception):
pass
logger.info("All SOCKS5 proxies stopped")
@property
def active_count(self) -> int:
return len(self._proxies)
+112
View File
@@ -0,0 +1,112 @@
"""Service for scheduled page visits."""
import asyncio
import random
import logging
from typing import Optional, Callable, Awaitable
from config.settings import settings
from services.browser_service import BrowserService
logger = logging.getLogger(__name__)
class VisitScheduler:
"""
Scheduler for periodic page visits.
Features:
- Schedule periodic visits
- Configurable delays and reading times
- Progress tracking
"""
def __init__(
self,
browser_service: BrowserService,
url: str,
min_delay: int = None,
max_delay: int = None,
min_reading: int = None,
max_reading: int = None,
max_visits: Optional[int] = None,
on_visit_complete: Optional[Callable] = None,
on_progress: Optional[Callable] = None
):
self.browser_service = browser_service
self.url = url
self.min_delay = min_delay or settings.DEFAULT_MIN_DELAY
self.max_delay = max_delay or settings.DEFAULT_MAX_DELAY
self.min_reading = min_reading or settings.DEFAULT_MIN_READING
self.max_reading = max_reading or settings.DEFAULT_MAX_READING
self.max_visits = max_visits
self.on_visit_complete = on_visit_complete
self.on_progress = on_progress
self.visit_count = 0
self.successful = 0
self.failed = 0
async def run(self) -> dict:
"""
Run scheduled visits.
Returns:
Statistics dict
"""
total_reading_time = 0
try:
while True:
if self.max_visits and self.visit_count >= self.max_visits:
logger.info(f"Visit limit reached: {self.max_visits}")
break
# Randomize parameters
current_reading = random.randint(self.min_reading, self.max_reading)
current_delay = random.randint(self.min_delay, self.max_delay)
self.visit_count += 1
logger.info(
f"Visit {self.visit_count}/{'' if not self.max_visits else self.max_visits}: "
f"reading={current_reading}s"
)
# Visit page
result = await self.browser_service.visit_page(
self.url, current_reading
)
if result.success:
self.successful += 1
total_reading_time += current_reading
else:
self.failed += 1
# Notify about visit completion
if self.on_visit_complete:
await self.on_visit_complete(
self.visit_count, result, current_delay, current_reading
)
if self.max_visits and self.visit_count >= self.max_visits:
break
# Wait before next visit
if self.on_progress:
await self.on_progress(self.visit_count, current_delay)
await asyncio.sleep(current_delay)
except asyncio.CancelledError:
logger.info("Visit scheduler cancelled")
raise
return {
"total_visits": self.visit_count,
"successful": self.successful,
"failed": self.failed,
"avg_reading_time": total_reading_time / self.successful if self.successful else 0
}