fix sock5 connect
This commit is contained in:
@@ -156,50 +156,53 @@ class Socks5ToHttpProxy:
|
|||||||
reader: asyncio.StreamReader,
|
reader: asyncio.StreamReader,
|
||||||
writer: asyncio.StreamWriter
|
writer: asyncio.StreamWriter
|
||||||
):
|
):
|
||||||
"""Обрабатывает CONNECT запрос."""
|
"""Обрабатывает CONNECT запрос через асинхронный SOCKS5."""
|
||||||
import socks
|
host = "unknown"
|
||||||
remote_socket = None
|
port = 443
|
||||||
|
|
||||||
try:
|
try:
|
||||||
parts = request_line.split()
|
parts = request_line.split()
|
||||||
if len(parts) < 2:
|
if len(parts) < 2:
|
||||||
writer.write(b'HTTP/1.1 400 Bad Request\r\n\r\n')
|
writer.write(b'HTTP/1.1 400 Bad Request\r\n\r\n')
|
||||||
return
|
return
|
||||||
|
|
||||||
target = parts[1]
|
target = parts[1]
|
||||||
if ':' in target:
|
if ':' in target:
|
||||||
host, port = target.rsplit(':', 1)
|
host, port = target.rsplit(':', 1)
|
||||||
port = int(port)
|
port = int(port)
|
||||||
else:
|
else:
|
||||||
host = target
|
host = target
|
||||||
port = 443
|
|
||||||
|
# Пропускаем заголовки
|
||||||
# Пропускаем заголовки (с таймаутом)
|
|
||||||
while True:
|
while True:
|
||||||
line = await asyncio.wait_for(reader.readline(), timeout=5)
|
line = await asyncio.wait_for(reader.readline(), timeout=5)
|
||||||
if not line or line in [b'\r\n', b'\n']:
|
if not line or line in [b'\r\n', b'\n']:
|
||||||
break
|
break
|
||||||
|
|
||||||
# Создаем SOCKS5 соединение
|
# Асинхронное SOCKS5 соединение через python-socks
|
||||||
remote_socket = socks.socksocket()
|
from python_socks.async_.asyncio import Proxy
|
||||||
remote_socket.settimeout(15)
|
from python_socks import ProxyType
|
||||||
|
|
||||||
if self.username:
|
if self.username:
|
||||||
logger.info(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=yes user={self.username[:4]}***")
|
logger.info(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=yes user={self.username[:4]}***")
|
||||||
remote_socket.set_proxy(
|
proxy = Proxy.create(
|
||||||
socks.SOCKS5, self.socks5_host, self.socks5_port,
|
ProxyType.SOCKS5,
|
||||||
username=self.username, password=self.password
|
self.socks5_host, self.socks5_port,
|
||||||
|
username=self.username, password=self.password,
|
||||||
|
rdns=True
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
logger.warning(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=NO — credentials missing!")
|
logger.warning(f"SOCKS5 connect {host}:{port} via {self.socks5_host}:{self.socks5_port} auth=NO")
|
||||||
remote_socket.set_proxy(
|
proxy = Proxy.create(
|
||||||
socks.SOCKS5, self.socks5_host, self.socks5_port
|
ProxyType.SOCKS5,
|
||||||
|
self.socks5_host, self.socks5_port,
|
||||||
|
rdns=True
|
||||||
)
|
)
|
||||||
|
|
||||||
# Подключаемся (в отдельном потоке чтобы не блокировать)
|
|
||||||
loop = asyncio.get_running_loop()
|
|
||||||
try:
|
try:
|
||||||
await loop.run_in_executor(None, remote_socket.connect, (host, port))
|
sock = await asyncio.wait_for(
|
||||||
|
proxy.connect(dest_host=host, dest_port=port),
|
||||||
|
timeout=15
|
||||||
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"SOCKS5 connect FAILED {host}:{port} via {self.socks5_host}: {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')
|
writer.write(b'HTTP/1.1 502 Bad Gateway\r\n\r\n')
|
||||||
@@ -207,27 +210,51 @@ class Socks5ToHttpProxy:
|
|||||||
return
|
return
|
||||||
|
|
||||||
logger.info(f"SOCKS5 connected {host}:{port} OK")
|
logger.info(f"SOCKS5 connected {host}:{port} OK")
|
||||||
# Отвечаем клиенту
|
|
||||||
writer.write(b'HTTP/1.1 200 Connection Established\r\n\r\n')
|
writer.write(b'HTTP/1.1 200 Connection Established\r\n\r\n')
|
||||||
await writer.drain()
|
await writer.drain()
|
||||||
|
|
||||||
# Пересылка данных с таймаутом
|
# Полностью асинхронный relay через asyncio streams
|
||||||
|
remote_reader, remote_writer = await asyncio.open_connection(sock=sock)
|
||||||
await asyncio.wait_for(
|
await asyncio.wait_for(
|
||||||
self._relay(reader, writer, remote_socket),
|
self._relay_streams(reader, writer, remote_reader, remote_writer),
|
||||||
timeout=60
|
timeout=120
|
||||||
)
|
)
|
||||||
|
|
||||||
except asyncio.TimeoutError:
|
except asyncio.TimeoutError:
|
||||||
logger.warning(f"CONNECT relay timeout {host}:{port}")
|
logger.warning(f"CONNECT relay timeout {host}:{port}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"CONNECT error {host}:{port}: {e}")
|
logger.warning(f"CONNECT error {host}:{port}: {e}")
|
||||||
finally:
|
|
||||||
if remote_socket:
|
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:
|
try:
|
||||||
remote_socket.close()
|
dst.close()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
|
await asyncio.gather(
|
||||||
|
pipe(client_reader, remote_writer),
|
||||||
|
pipe(remote_reader, client_writer),
|
||||||
|
return_exceptions=True
|
||||||
|
)
|
||||||
|
|
||||||
async def _relay(
|
async def _relay(
|
||||||
self,
|
self,
|
||||||
client_reader: asyncio.StreamReader,
|
client_reader: asyncio.StreamReader,
|
||||||
|
|||||||
Reference in New Issue
Block a user