Files
pleiades/gateway.py

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