427 lines
15 KiB
Python
427 lines
15 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.debug(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 запрос через асинхронный SOCKS5."""
|
|
host = "unknown"
|
|
port = 443
|
|
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
|
|
|
|
# Пропускаем заголовки
|
|
while True:
|
|
line = await asyncio.wait_for(reader.readline(), timeout=5)
|
|
if not line or line in [b'\r\n', b'\n']:
|
|
break
|
|
|
|
# Асинхронное SOCKS5 соединение через python-socks
|
|
from python_socks.async_.asyncio import Proxy
|
|
from python_socks import ProxyType
|
|
|
|
if self.username:
|
|
logger.debug(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=yes")
|
|
proxy = Proxy.create(
|
|
ProxyType.SOCKS5,
|
|
self.socks5_host, self.socks5_port,
|
|
username=self.username, password=self.password,
|
|
rdns=True
|
|
)
|
|
else:
|
|
logger.warning(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=NO")
|
|
proxy = Proxy.create(
|
|
ProxyType.SOCKS5,
|
|
self.socks5_host, self.socks5_port,
|
|
rdns=True
|
|
)
|
|
|
|
try:
|
|
sock = await asyncio.wait_for(
|
|
proxy.connect(dest_host=host, dest_port=port),
|
|
timeout=15
|
|
)
|
|
except Exception as e:
|
|
logger.warning(f"SOCKS5 connect FAILED {host}:{port} via {self.socks5_host}: {e}")
|
|
writer.write(b'HTTP/1.1 502 Bad Gateway\r\n\r\n')
|
|
await writer.drain()
|
|
return
|
|
|
|
logger.debug(f"SOCKS5 connected {host}:{port} OK")
|
|
writer.write(b'HTTP/1.1 200 Connection Established\r\n\r\n')
|
|
await writer.drain()
|
|
|
|
# Полностью асинхронный relay через asyncio streams
|
|
remote_reader, remote_writer = await asyncio.open_connection(sock=sock)
|
|
await asyncio.wait_for(
|
|
self._relay_streams(reader, writer, remote_reader, remote_writer),
|
|
timeout=120
|
|
)
|
|
|
|
except asyncio.TimeoutError:
|
|
logger.warning(f"CONNECT relay timeout {host}:{port}")
|
|
except Exception as e:
|
|
logger.warning(f"CONNECT error {host}:{port}: {e}")
|
|
|
|
async def _relay_streams(
|
|
self,
|
|
client_reader: asyncio.StreamReader,
|
|
client_writer: asyncio.StreamWriter,
|
|
remote_reader: asyncio.StreamReader,
|
|
remote_writer: asyncio.StreamWriter,
|
|
):
|
|
"""Асинхронная двунаправленная пересылка через asyncio streams."""
|
|
async def pipe(src: asyncio.StreamReader, dst: asyncio.StreamWriter):
|
|
try:
|
|
while True:
|
|
data = await src.read(65536)
|
|
if not data:
|
|
break
|
|
dst.write(data)
|
|
await dst.drain()
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
try:
|
|
dst.close()
|
|
except Exception:
|
|
pass
|
|
|
|
await asyncio.gather(
|
|
pipe(client_reader, remote_writer),
|
|
pipe(remote_reader, client_writer),
|
|
return_exceptions=True
|
|
)
|
|
|
|
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
|
|
|
|
# Быстрая проверка под lock'ом
|
|
async with self._lock:
|
|
if key in self._proxies:
|
|
self._last_used[key] = asyncio.get_running_loop().time()
|
|
return self._proxies[key].proxy_config_for_browser
|
|
|
|
# proxy.start() — медленная операция, запускаем вне lock'а
|
|
proxy = Socks5ToHttpProxy(
|
|
socks5_host=socks5_proxy.ip,
|
|
socks5_port=int(socks5_proxy.port),
|
|
username=socks5_proxy.login,
|
|
password=socks5_proxy.password
|
|
)
|
|
await proxy.start()
|
|
|
|
# Добавляем под lock'ом (двойная проверка — другая задача могла успеть)
|
|
async with self._lock:
|
|
if key not in self._proxies:
|
|
self._proxies[key] = proxy
|
|
else:
|
|
await proxy.stop() # уже создан другой задачей
|
|
self._last_used[key] = asyncio.get_running_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) |