244 lines
8.4 KiB
Python
244 lines
8.4 KiB
Python
#!/usr/bin/env python3
|
|
"""Pleiades gateway — SOCKS5 + HTTP CONNECT proxy.
|
|
|
|
Accepts customer user:pass, looks up their geo prefs, rewrites the upstream
|
|
IPRoyal password with geo/sticky suffixes, and relays. Meters bytes per user
|
|
for GB billing. Does NOT log traffic content — only connection metadata.
|
|
|
|
Billing note: bytes are metered (required for GB plans). Target hostnames and
|
|
payload are never persisted.
|
|
"""
|
|
import asyncio
|
|
import base64
|
|
import logging
|
|
import signal
|
|
import struct
|
|
|
|
import db
|
|
|
|
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
|
|
log = logging.getLogger("pleiades-gateway")
|
|
|
|
UP = db.CFG["upstream"]
|
|
GW = db.CFG["gateway"]
|
|
UP_AUTH = base64.b64encode(f"{UP['username']}:{UP['password']}".encode()).decode()
|
|
|
|
|
|
def _sticky_sid(user):
|
|
"""Stable 8-char alphanumeric session id per user (sticky IP for lifetime window)."""
|
|
import hashlib
|
|
return hashlib.md5(("u%d" % user["id"]).encode()).hexdigest()[:8]
|
|
|
|
|
|
def build_upstream_password(user):
|
|
"""Append geo/sticky suffixes to the base upstream password."""
|
|
p = UP["password"]
|
|
if user.get("region"):
|
|
p += f"_region-{user['region']}"
|
|
if user.get("country"):
|
|
p += f"_country-{user['country']}"
|
|
if user.get("city"):
|
|
p += f"_city-{user['city']}"
|
|
if user.get("sticky"):
|
|
p += f"_session-{_sticky_sid(user)}_lifetime-30m"
|
|
return p
|
|
|
|
|
|
async def connect_upstream(host, port, user):
|
|
"""Open TCP to IPRoyal, issue HTTP CONNECT, return (reader, writer, leftover)."""
|
|
up_pass = build_upstream_password(user)
|
|
auth = base64.b64encode(f"{UP['username']}:{up_pass}".encode()).decode()
|
|
reader, writer = await asyncio.open_connection(UP["host"], UP["port"])
|
|
req = (f"CONNECT {host}:{port} HTTP/1.1\r\n"
|
|
f"Host: {host}:{port}\r\n"
|
|
f"Proxy-Authorization: Basic {auth}\r\n\r\n")
|
|
writer.write(req.encode())
|
|
await writer.drain()
|
|
data = b""
|
|
while b"\r\n\r\n" not in data:
|
|
chunk = await reader.read(4096)
|
|
if not chunk:
|
|
break
|
|
data += chunk
|
|
if len(data) > 65536:
|
|
break
|
|
head, _, leftover = data.partition(b"\r\n\r\n")
|
|
status = head.decode(errors="replace").split("\r\n")[0]
|
|
if "200" not in status:
|
|
raise ConnectionError(f"upstream refused: {status}")
|
|
return reader, writer, leftover
|
|
|
|
|
|
async def pipe(reader, writer, counter):
|
|
try:
|
|
while True:
|
|
data = await reader.read(65536)
|
|
if not data:
|
|
break
|
|
counter[0] += len(data)
|
|
writer.write(data)
|
|
await writer.drain()
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
try:
|
|
writer.close()
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def parse_target(reader, data):
|
|
"""Parse SOCKS5 address from accumulated bytes. Returns (host, port, consumed)."""
|
|
atyp = data[3]
|
|
if atyp == 0x01: # IPv4
|
|
host = ".".join(str(b) for b in data[4:8])
|
|
port = struct.unpack(">H", data[8:10])[0]
|
|
return host, port, 10
|
|
if atyp == 0x03: # domain
|
|
ln = data[4]
|
|
host = data[5:5 + ln].decode(errors="replace")
|
|
port = struct.unpack(">H", data[5 + ln:7 + ln])[0]
|
|
return host, port, 7 + ln
|
|
if atyp == 0x04: # IPv6
|
|
import ipaddress
|
|
host = str(ipaddress.IPv6Address(data[4:20]))
|
|
port = struct.unpack(">H", data[20:22])[0]
|
|
return host, port, 22
|
|
raise ValueError("bad atyp")
|
|
|
|
|
|
async def handle_socks(reader, writer):
|
|
peername = writer.get_extra_info("peername")
|
|
user = None
|
|
counter = [0]
|
|
try:
|
|
# greeting
|
|
head = await reader.readexactly(2)
|
|
ver, nmethods = head[0], head[1]
|
|
methods = await reader.readexactly(nmethods)
|
|
if ver != 5 or 0x02 not in methods:
|
|
writer.write(b"\x05\xff")
|
|
await writer.drain()
|
|
return
|
|
writer.write(b"\x05\x02") # choose user/pass auth
|
|
await writer.drain()
|
|
# auth (RFC 1929)
|
|
await reader.readexactly(1) # version 0x01
|
|
ulen = (await reader.readexactly(1))[0]
|
|
username = (await reader.readexactly(ulen)).decode()
|
|
plen = (await reader.readexactly(1))[0]
|
|
password = (await reader.readexactly(plen)).decode()
|
|
user = db.lookup_user_by_proxy(username, password)
|
|
if user is None:
|
|
writer.write(b"\x01\x01")
|
|
await writer.drain()
|
|
return
|
|
writer.write(b"\x01\x00")
|
|
await writer.drain()
|
|
# request header: ver, cmd, rsv, atyp
|
|
hdr = await reader.readexactly(4)
|
|
cmd = hdr[1]
|
|
atyp = hdr[3]
|
|
if cmd != 0x01: # only CONNECT supported
|
|
writer.write(b"\x05\x07\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
await writer.drain()
|
|
return
|
|
if atyp == 0x01: # IPv4
|
|
host = ".".join(str(b) for b in (await reader.readexactly(4)))
|
|
elif atyp == 0x03: # domain
|
|
ln = (await reader.readexactly(1))[0]
|
|
host = (await reader.readexactly(ln)).decode(errors="replace")
|
|
elif atyp == 0x04: # IPv6
|
|
import ipaddress
|
|
host = str(ipaddress.IPv6Address(await reader.readexactly(16)))
|
|
else:
|
|
raise ValueError("bad atyp")
|
|
port = struct.unpack(">H", await reader.readexactly(2))[0]
|
|
# connect upstream
|
|
up_reader, up_writer, leftover = await connect_upstream(host, port, user)
|
|
# success reply
|
|
writer.write(b"\x05\x00\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
await writer.drain()
|
|
# forward any leftover bytes from upstream to client
|
|
if leftover:
|
|
counter[0] += len(leftover)
|
|
writer.write(leftover)
|
|
await writer.drain()
|
|
# relay both directions
|
|
await asyncio.gather(
|
|
pipe(reader, up_writer, counter),
|
|
pipe(up_reader, writer, counter),
|
|
)
|
|
except (asyncio.IncompleteReadError, ConnectionError, ValueError, struct.error) as e:
|
|
log.debug("socks %s closed: %s", peername, e)
|
|
except Exception as e:
|
|
log.warning("socks %s error: %s", peername, e)
|
|
finally:
|
|
if user:
|
|
db.add_usage(user["id"], counter[0])
|
|
|
|
|
|
async def handle_http(reader, writer):
|
|
peername = writer.get_extra_info("peername")
|
|
user = None
|
|
counter = [0]
|
|
try:
|
|
first = await reader.readuntil(b"\r\n")
|
|
parts = first.decode(errors="replace").split()
|
|
if len(parts) < 2 or parts[0].upper() != "CONNECT":
|
|
writer.write(b"HTTP/1.1 405 Method Not Allowed\r\n\r\n")
|
|
await writer.drain()
|
|
return
|
|
target = parts[1]
|
|
host, _, port = target.rpartition(":")
|
|
# read headers to find proxy auth
|
|
auth = None
|
|
while True:
|
|
line = await reader.readuntil(b"\r\n")
|
|
if line == b"\r\n":
|
|
break
|
|
if line.lower().startswith(b"proxy-authorization:"):
|
|
auth = line.decode().split(":", 1)[1].strip()
|
|
if auth and auth.lower().startswith("basic "):
|
|
raw = base64.b64decode(auth[6:]).decode()
|
|
username, _, password = raw.partition(":")
|
|
user = db.lookup_user_by_proxy(username, password)
|
|
if user is None:
|
|
writer.write(b"HTTP/1.1 407 Proxy Authentication Required\r\n\r\n")
|
|
await writer.drain()
|
|
return
|
|
up_reader, up_writer, leftover = await connect_upstream(host, int(port), user)
|
|
writer.write(b"HTTP/1.1 200 Connection Established\r\n\r\n")
|
|
await writer.drain()
|
|
if leftover:
|
|
counter[0] += len(leftover)
|
|
writer.write(leftover)
|
|
await writer.drain()
|
|
await asyncio.gather(
|
|
pipe(reader, up_writer, counter),
|
|
pipe(up_reader, writer, counter),
|
|
)
|
|
except (asyncio.IncompleteReadError, ConnectionError, ValueError) as e:
|
|
log.debug("http %s closed: %s", peername, e)
|
|
except Exception as e:
|
|
log.warning("http %s error: %s", peername, e)
|
|
finally:
|
|
if user:
|
|
db.add_usage(user["id"], counter[0])
|
|
|
|
|
|
async def main():
|
|
db.init_db()
|
|
socks = await asyncio.start_server(handle_socks, GW["listen"], GW["socks_port"])
|
|
http = await asyncio.start_server(handle_http, GW["listen"], GW["http_port"])
|
|
log.info("SOCKS5 on :%d HTTP on :%d", GW["socks_port"], GW["http_port"])
|
|
async with socks, http:
|
|
await asyncio.Future()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
try:
|
|
asyncio.run(main())
|
|
except KeyboardInterrupt:
|
|
pass
|