Files
Click/services/socks5_to_http_proxy.py
T
Yuriy Yuriev 46c325540b fix logs
2026-05-19 13:41:59 +07:00

384 lines
14 KiB
Python

"""
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()
logger.info(f"Tunnel request: {request_text[:80]}")
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.info(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} 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}:{self.socks5_port} 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)