""" 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: logger.debug(f"SOCKS5 connect {host}:{port} via {self.socks5_host} auth=yes user={self.username[:4]}***") remote_socket.set_proxy( socks.SOCKS5, self.socks5_host, self.socks5_port, username=self.username, password=self.password ) else: logger.warning(f"SOCKS5 connect {host}:{port} via {self.socks5_host} auth=NO — credentials missing!") 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)