diff --git a/app.py b/app.py index 2a6e7e0..e15cfbf 100644 --- a/app.py +++ b/app.py @@ -38,6 +38,10 @@ BTCPAY_URL = os.environ.get("BTCPAY_URL", "https://10.30.20.140") BTCPAY_STORE = os.environ.get("BTCPAY_STORE_ID", "") BTCPAY_TOKEN = os.environ.get("BTCPAY_API_KEY", "") +# Public SOCKS5 endpoint customers connect to (authenticated frontend on CT158) +PROXY_HOST = os.environ.get("PROXY_PUBLIC_HOST", "10.30.20.116") +PROXY_PORT = os.environ.get("PROXY_PUBLIC_PORT", "1081") + def db(): con = getattr(g, "_db", None) if con is None: @@ -108,7 +112,8 @@ def index(): user = current_user() if not user: return render_template_string(AUTH_HTML, mode="login", error="") - return render_template_string(DASH_HTML, user=user, locations=LOCATIONS, plans=PLANS) + return render_template_string(DASH_HTML, user=user, locations=LOCATIONS, plans=PLANS, + proxy_host=PROXY_HOST, proxy_port=PROXY_PORT) @app.route("/health") def health(): @@ -255,12 +260,12 @@ def btcpay_webhook(): return jsonify({"ok": True}) def provision_proxy(location, user, pw): - """Best-effort: add user auth to the gost frontend for this location. - The gost SOCKS5 frontends need per-user auth; for now we record the creds - and the actual gost auth layer is wired per-location (see rigel-proxy setup).""" - # Placeholder — the real gost auth is applied via the proxy gateway config. - # Creds are persisted on the subscription row so the customer sees them. - app.logger.info(f"provision proxy {location} for {user}") + """No per-node writes needed: the authenticated SOCKS5 frontend + (rigel-proxy.service on CT158) validates every connection live against the + subscriptions table, so persisting the row IS the provisioning step. + The upstream exits are dumb unauthenticated relays on the LAN.""" + up = LOCATIONS.get(location, {}).get("upstream", "?") + app.logger.info(f"provisioned {location} (upstream {up}) for {user}") # ── templates ─────────────────────────────────────────────────────────── AUTH_HTML = r""" @@ -332,6 +337,10 @@ main{padding:28px;max-width:900px;margin:0 auto} .subs{margin-top:32px}.subs h2{font-size:17px}.subs table{width:100%;border-collapse:collapse;font-size:13px} .subs th,.subs td{padding:10px;text-align:left;border-bottom:1px solid #1e2742}.subs th{color:var(--muted);font-weight:500} .badge{padding:3px 8px;border-radius:10px;font-size:11px}.active{background:#14321f;color:#2ecc71}.pending{background:#332a14;color:#f1c40f}.expired{background:#33171a;color:#e74c3c} +.subcard{background:var(--card);border:1px solid #1e2742;border-radius:12px;padding:18px;margin-top:14px} +.srow{font-size:15px;margin-bottom:6px} +.lbl{font-size:11px;color:var(--muted);text-transform:uppercase;letter-spacing:.5px;margin-top:10px} +code{display:block;background:#0d1325;border:1px solid #26304f;border-radius:6px;padding:8px 10px;font-size:12px;color:#9fd0ff;word-break:break-all;margin-top:3px;font-family:ui-monospace,SFMono-Regular,Menlo,monospace}

✦ Rigel — {{user["username"]}}

@@ -359,9 +368,21 @@ async function logout(){await fetch('/api/logout',{method:'POST'});location.relo async function loadSubs(){ let r=await fetch('/api/subscriptions').then(x=>x.json()); if(!r.length){document.getElementById('subs').innerHTML='

No subscriptions yet.

';return;} - let html=''; - for(const s of r){html+=``;} - html+='
LocationPlanStatusExpiresProxy creds
${s.location}${s.plan}${s.status}${s.expires_at||'—'}${s.proxy_user||'—'}
';document.getElementById('subs').innerHTML=html; + let h=''; + for(const s of r){ + h+=`
${s.location} · ${s.plan} · ${s.status}
`; + if(s.status==='active'&&s.proxy_user){ + h+=`
SOCKS5 endpoint
{{proxy_host}}:{{proxy_port}}` + + `
Username
${s.proxy_user}` + + `
Password
${s.proxy_pass||'\u2014'}` + + `
Expires (UTC)
${s.expires_at||'\u2014'}` + + `
Test it
curl --socks5-hostname ${s.proxy_user}:${s.proxy_pass}@{{proxy_host}}:{{proxy_port}} https://api.ipify.org`; + } else { + h+=`
Awaiting payment confirmation\u2026 credentials appear here automatically.
`; + } + h+='
'; + } + document.getElementById('subs').innerHTML=h; } loadSubs(); diff --git a/proxy_server.py b/proxy_server.py new file mode 100644 index 0000000..dd9f079 --- /dev/null +++ b/proxy_server.py @@ -0,0 +1,264 @@ +#!/usr/bin/env python3 +""" +Rigel — authenticated SOCKS5 frontend (customer entry point). + +Customers authenticate with the per-user credentials issued when their +invoice settles (subscriptions.proxy_user / proxy_pass). We validate LIVE +against rigel.db, resolve the subscription's location, and relay to that +location's upstream SOCKS5 exit on the LAN. + +Upstreams are derived from app.LOCATIONS so the two cannot drift. Upstreams +that themselves need auth (IPRoyal residential/mobile) are supported via +UPSTREAM__USER / UPSTREAM__PASS env vars. + +Implements SOCKS5 RFC1928 + username/password auth RFC1929. +""" +import asyncio +import importlib.util +import ipaddress +import logging +import os +import sqlite3 +import struct +from datetime import datetime + +logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)s %(message)s") +log = logging.getLogger("rigel-proxy") + +HERE = os.path.dirname(os.path.abspath(__file__)) +DB = os.environ.get("RIGEL_DB", os.path.join(HERE, "rigel.db")) +LISTEN_HOST = os.environ.get("PROXY_LISTEN_HOST", "0.0.0.0") +LISTEN_PORT = int(os.environ.get("PROXY_LISTEN_PORT", "1081")) + +DEFAULT_LOCATIONS = { + "tokyo": {"upstream": "10.30.20.154:1080"}, + "london": {"upstream": "10.30.20.71:1080"}, + "sydney": {"upstream": "10.30.20.189:1080"}, +} + + +def _load_upstreams(): + """Derive upstreams from app.LOCATIONS (single source of truth).""" + locs = DEFAULT_LOCATIONS + try: + spec = importlib.util.spec_from_file_location( + "rigel_app", os.path.join(HERE, "app.py")) + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + locs = mod.LOCATIONS + log.info("upstreams loaded from app.LOCATIONS") + except Exception as exc: # keep serving on built-ins rather than dying + log.warning("could not import app.LOCATIONS (%s); using defaults", exc) + out = {} + for key, val in locs.items(): + up = (val or {}).get("upstream", "") + if ":" not in up: + continue + host, port = up.rsplit(":", 1) + try: + port = int(port) + except ValueError: + continue + out[key] = { + "host": host, + "port": port, + "user": os.environ.get(f"UPSTREAM_{key.upper()}_USER") or None, + "pass": os.environ.get(f"UPSTREAM_{key.upper()}_PASS") or None, + } + return out + + +UPSTREAMS = _load_upstreams() + + +def authenticate(user, pw): + """Validate per-user creds live against the DB. -> (location, upstream)|None.""" + try: + con = sqlite3.connect(DB, timeout=5) + except sqlite3.Error as exc: + log.error("db open failed: %s", exc) + return None + con.row_factory = sqlite3.Row + try: + row = con.execute( + "SELECT location, proxy_pass, status, expires_at FROM subscriptions " + "WHERE proxy_user=? ORDER BY id DESC LIMIT 1", (user,) + ).fetchone() + except sqlite3.Error as exc: + log.error("db query failed: %s", exc) + return None + finally: + con.close() + if not row or row["proxy_pass"] != pw: + return None + if row["status"] != "active": + return None + exp = row["expires_at"] + if exp: + try: + # webhook stores naive UTC (datetime.utcnow().isoformat()) + if datetime.fromisoformat(exp) < datetime.utcnow(): + log.info("expired sub user=%s exp=%s", user, exp) + return None + except ValueError: + pass + up = UPSTREAMS.get(row["location"]) + if not up: + return None + return row["location"], up + + +async def _upstream_open(up, atyp, addr_bytes, port): + """Open a SOCKS5 connection to an upstream exit for the client's target.""" + ur, uw = await asyncio.wait_for( + asyncio.open_connection(up["host"], up["port"]), timeout=20) + if up.get("user"): + uw.write(b"\x05\x02\x00\x02") + await uw.drain() + if await ur.readexactly(2) != b"\x05\x02": + raise IOError("upstream refused user/pass auth") + u = up["user"].encode() + p = (up["pass"] or "").encode() + uw.write(b"\x01" + bytes([len(u)]) + u + bytes([len(p)]) + p) + await uw.drain() + if (await ur.readexactly(2))[1] != 0: + raise IOError("upstream auth rejected") + else: + uw.write(b"\x05\x01\x00") + await uw.drain() + if await ur.readexactly(2) != b"\x05\x00": + raise IOError("upstream refused no-auth") + # NOTE: for ATYP=3 the domain MUST be length-prefixed; without it the + # upstream reads the first domain byte ('a' = 0x61 = 97) as the length + # and blocks forever waiting for a 97-byte hostname. + addr_field = bytes([len(addr_bytes)]) + addr_bytes if atyp == 3 else addr_bytes + uw.write(b"\x05\x01\x00" + bytes([atyp]) + addr_field + struct.pack(">H", port)) + await uw.drain() + rep = await ur.readexactly(4) + if rep[1] != 0: + raise IOError(f"upstream CONNECT rep={rep[1]}") + if rep[3] == 1: + await ur.readexactly(4) + elif rep[3] == 4: + await ur.readexactly(16) + elif rep[3] == 3: + await ur.readexactly((await ur.readexactly(1))[0]) + await ur.readexactly(2) + return ur, uw + + +async def _pipe(reader, writer): + try: + while True: + data = await reader.read(65536) + if not data: + break + writer.write(data) + await writer.drain() + except (OSError, asyncio.IncompleteReadError): + pass + finally: + try: + writer.close() + except OSError: + pass + + +def _deny(writer, rep=0x01): + writer.write(b"\x05" + bytes([rep]) + b"\x00\x01" + b"\x00" * 4 + b"\x00\x00") + + +async def handle(reader, writer): + peer = writer.get_extra_info("peername") + ip = peer[0] if peer else "?" + try: + ver, nmeth = await reader.readexactly(2) + if ver != 5: + return + methods = await reader.readexactly(nmeth) + if 0x02 not in methods: + writer.write(b"\x05\xff") + await writer.drain() + return + writer.write(b"\x05\x02") + await writer.drain() + + if (await reader.readexactly(1))[0] != 1: + return + ulen = (await reader.readexactly(1))[0] + uname = (await reader.readexactly(ulen)).decode(errors="replace") + plen = (await reader.readexactly(1))[0] + passwd = (await reader.readexactly(plen)).decode(errors="replace") + + auth = authenticate(uname, passwd) + if not auth: + log.info("AUTH FAIL user=%r from %s", uname, ip) + writer.write(b"\x01\x01") + await writer.drain() + return + location, up = auth + writer.write(b"\x01\x00") + await writer.drain() + + req = await reader.readexactly(4) + cmd, atyp = req[1], req[3] + if cmd != 1: + _deny(writer, 0x07) # command not supported + await writer.drain() + return + if atyp == 1: + addr_bytes = await reader.readexactly(4) + target = str(ipaddress.IPv4Address(addr_bytes)) + elif atyp == 3: + n = (await reader.readexactly(1))[0] + addr_bytes = await reader.readexactly(n) + target = addr_bytes.decode(errors="replace") + elif atyp == 4: + addr_bytes = await reader.readexactly(16) + target = str(ipaddress.IPv6Address(addr_bytes)) + else: + _deny(writer, 0x08) + await writer.drain() + return + port = struct.unpack(">H", await reader.readexactly(2))[0] + + try: + ur, uw = await _upstream_open(up, atyp, addr_bytes, port) + except (OSError, asyncio.IncompleteReadError, asyncio.TimeoutError, IOError) as exc: + log.warning("UPSTREAM FAIL user=%r loc=%s target=%s:%s err=%s", + uname, location, target, port, exc) + _deny(writer, 0x04) # host unreachable + await writer.drain() + return + + writer.write(b"\x05\x00\x00\x01" + b"\x00" * 4 + b"\x00\x00") + await writer.drain() + log.info("OK user=%r loc=%s -> %s:%s from %s", uname, location, target, port, ip) + + await asyncio.gather(_pipe(reader, uw), _pipe(ur, writer)) + except (asyncio.IncompleteReadError, ConnectionResetError, BrokenPipeError): + pass + except Exception as exc: + log.exception("handler error from %s: %s", ip, exc) + finally: + try: + writer.close() + except OSError: + pass + + +async def main(): + server = await asyncio.start_server(handle, LISTEN_HOST, LISTEN_PORT) + addrs = ", ".join(str(s.getsockname()) for s in server.sockets) + log.info("rigel proxy listening on %s | locations=%s | db=%s", + addrs, sorted(UPSTREAMS), DB) + async with server: + await server.serve_forever() + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + pass