#!/usr/bin/env python3
"""
PCA Phobos — Web Panel for Phobos (Obfuscated WireGuard VPN)
Management panel: clients, sessions, labels, subscriptions, Telegram alerts.
"""
import json, os, subprocess, threading, time, secrets, hashlib, re, fcntl
from datetime import datetime, timedelta
from pathlib import Path
from flask import Flask, request, redirect, url_for, session, make_response
_settings_lock = threading.Lock()
app = Flask(__name__)
PHOBOS_DIR = "/opt/Phobos"
CLIENTS_DIR = f"{PHOBOS_DIR}/clients"
SERVER_ENV = f"{PHOBOS_DIR}/server/server.env"
PANEL_DIR = "/opt/phobos-panel"
SETTINGS_FILE = f"{PANEL_DIR}/settings.json"
SECRET_FILE = f"{PANEL_DIR}/.secret_key"
SERVER_IP = subprocess.getoutput("curl -s https://api.ipify.org 2>/dev/null || hostname -I | awk '{print $1}'").strip()
os.makedirs(PANEL_DIR, exist_ok=True)
if os.path.exists(SECRET_FILE):
app.secret_key = open(SECRET_FILE).read().strip()
else:
app.secret_key = secrets.token_hex(32)
with open(SECRET_FILE, "w") as f:
f.write(app.secret_key)
DEFAULT_SETTINGS = {
"admin_pass": "OcAdmin2026!",
"tg_bot_token": "",
"tg_chat_id": "",
"monitor_interval": 30,
"labels": {},
"subscriptions": {}
}
def load_settings():
with _settings_lock:
if os.path.exists(SETTINGS_FILE):
with open(SETTINGS_FILE) as f:
s = json.load(f)
for k, v in DEFAULT_SETTINGS.items():
s.setdefault(k, v)
return s
return dict(DEFAULT_SETTINGS)
def save_settings(s):
with _settings_lock:
with open(SETTINGS_FILE, "w") as f:
json.dump(s, f, indent=2, ensure_ascii=False)
def tg_send(text):
s = load_settings()
token, chat = s.get("tg_bot_token", ""), s.get("tg_chat_id", "")
if not token or not chat:
return
try:
import urllib.request
url = f"https://api.telegram.org/bot{token}/sendMessage"
data = json.dumps({"chat_id": chat, "text": text, "parse_mode": "HTML"}).encode()
req = urllib.request.Request(url, data=data, headers={"Content-Type": "application/json"})
urllib.request.urlopen(req, timeout=10)
except Exception:
pass
def notify_on(cid):
return load_settings().get("client_notify", {}).get(cid, True)
def create_client_cli(name):
"""Create a client (add + assign + fanout + install link). Returns (ok, html_msg). Used by web + Telegram bot."""
import re as _re
name = (name or "").strip()
if not name or not _re.match(r"^[a-zA-Z0-9_-]+$", name):
return False, "Имя: латиница, цифры, _ и - (без пробелов)."
if os.path.isdir(os.path.join(CLIENTS_DIR, name)):
return False, "Клиент %s уже существует." % name
script = "%s/repo/server/scripts/phobos-client.sh" % PHOBOS_DIR
try:
subprocess.check_output([script, "add", name], text=True, timeout=45,
stderr=subprocess.STDOUT, env={**os.environ, **_load_server_env()})
except subprocess.CalledProcessError as e:
return False, "Ошибка создания: %s" % (e.output or "")[-300:]
try:
auto_assign_server(name); fanout_router_config(name)
except Exception:
pass
cmd = ""
try:
out = subprocess.check_output([script, "link", name], text=True, timeout=25,
stderr=subprocess.STDOUT, env={**os.environ, **_load_server_env()})
for ln in out.splitlines():
if ln.strip().startswith("curl") or "/init/" in ln:
cmd = ln.strip(); break
except Exception:
pass
msg = "\u2705 Клиент %s создан." % name
if cmd:
msg += "\n\nКоманда установки на роутер:\n%s" % cmd
return True, msg
def delete_client_cli(cid):
cid = (cid or "").strip()
if not cid or not os.path.isdir(os.path.join(CLIENTS_DIR, cid)):
return False, "Клиент %s не найден." % cid
script = "%s/repo/server/scripts/phobos-client.sh" % PHOBOS_DIR
try:
subprocess.check_output([script, "remove", cid], text=True, timeout=25,
stderr=subprocess.STDOUT, env={**os.environ, **_load_server_env()})
except subprocess.CalledProcessError as e:
return False, "Ошибка удаления: %s" % (e.output or "")[-300:]
try:
st = load_settings()
st.get("client_assignments", {}).pop(cid, None)
st.get("client_notify", {}).pop(cid, None)
save_settings(st)
except Exception:
pass
return True, "\U0001F5D1 Клиент %s удалён." % cid
def tg_api(method, payload, token=None):
import urllib.request
if token is None:
token = load_settings().get("tg_bot_token", "")
if not token:
return None
try:
data = json.dumps(payload).encode()
req = urllib.request.Request("https://api.telegram.org/bot%s/%s" % (token, method),
data=data, headers={"Content-Type": "application/json"})
return json.loads(urllib.request.urlopen(req, timeout=15).read())
except Exception:
return None
MAIN_KB = {
"keyboard": [
["➕ Создать клиента", "\U0001F4CB Список"],
["\U0001F5D1 Удалить", "\U0001F4CA Статус"],
["❓ Помощь"],
],
"resize_keyboard": True,
}
_bot_state = {}
def _bot_status_text():
servers = [{"ip": SERVER_IP, "is_primary": True}] + load_servers()
lines = ["\U0001F4CA Серверы:"]
for sv in servers:
ip = sv.get("ip", "")
st = "ok" if ip == SERVER_IP else server_stats_cache.get(ip, {}).get("status", "?")
emoji = "\U0001F7E2" if st == "ok" else ("\U0001F534" if st == "down" else "⚪")
lines.append("%s %s%s" % (emoji, ip, " (primary)" if sv.get("is_primary") else ""))
online = set()
try:
for p, info in get_wg_peers().items():
if is_peer_online(info.get("latest handshake", "")):
online.add(p)
for _ip, hs in server_handshakes_cache.items():
for p, tsv in hs.items():
if tsv and (time.time() - tsv) < 180:
online.add(p)
except Exception:
pass
lines.append("\n\U0001F4F1 Клиенты:")
cl = get_clients()
if not cl:
lines.append("нет")
for c in cl:
on = c.get("public_key", "") in online
lines.append(("\U0001F7E2 " if on else "⚪ ") + c.get("client_id", ""))
return "\n".join(lines)
def qr_png_bytes(data):
import qrcode, io
buf = io.BytesIO(); qrcode.make(data).save(buf, format="PNG"); return buf.getvalue()
def tg_send_photo(chat, png, caption=""):
import urllib.request, secrets as _sec
token = load_settings().get("tg_bot_token", "")
if not token:
return None
boundary = "----phobos" + _sec.token_hex(8)
def fld(name, val):
return ("--%s\r\nContent-Disposition: form-data; name=\"%s\"\r\n\r\n%s\r\n" % (boundary, name, val)).encode()
body = fld("chat_id", str(chat)) + fld("caption", caption) + fld("parse_mode", "HTML")
body += ("--%s\r\nContent-Disposition: form-data; name=\"photo\"; filename=\"qr.png\"\r\nContent-Type: image/png\r\n\r\n" % boundary).encode()
body += png + ("\r\n--%s--\r\n" % boundary).encode()
try:
req = urllib.request.Request("https://api.telegram.org/bot%s/sendPhoto" % token, data=body,
headers={"Content-Type": "multipart/form-data; boundary=%s" % boundary})
return urllib.request.urlopen(req, timeout=25).read()
except Exception:
return None
def bot_send_new_client(chat, name):
ok, rep = create_client_cli(name)
tg_api("sendMessage", {"chat_id": chat, "text": rep, "parse_mode": "HTML", "reply_markup": MAIN_KB})
if ok:
try:
conf, link = build_phone_config(name, "android")
if link:
tg_send_photo(chat, qr_png_bytes(link), "📱 Android (PhobosWG) — отсканируй QR в приложении")
tg_api("sendMessage", {"chat_id": chat, "text": "phobos:// ссылка (Android):\n%s" % link, "parse_mode": "HTML"})
except Exception:
pass
def telegram_bot():
offset = 0
import urllib.request
while True:
try:
s = load_settings()
token = s.get("tg_bot_token", "")
chat = str(s.get("tg_chat_id", ""))
if not token or not chat:
time.sleep(15); continue
url = "https://api.telegram.org/bot%s/getUpdates?timeout=30&offset=%d" % (token, offset)
data = json.loads(urllib.request.urlopen(url, timeout=40).read())
for upd in data.get("result", []):
offset = upd["update_id"] + 1
cq = upd.get("callback_query")
if cq:
if str(cq.get("message", {}).get("chat", {}).get("id", "")) != chat:
continue
cdata = cq.get("data", "")
if cdata.startswith("del:"):
_, rep = delete_client_cli(cdata[4:])
tg_api("answerCallbackQuery", {"callback_query_id": cq["id"], "text": "OK"})
tg_api("sendMessage", {"chat_id": chat, "text": rep, "parse_mode": "HTML", "reply_markup": MAIN_KB})
continue
m = upd.get("message") or upd.get("channel_post")
if not m:
continue
if str(m.get("chat", {}).get("id", "")) != chat:
continue
text = (m.get("text") or "").strip()
if not text:
continue
# detect reply-keyboard button by its emoji prefix (NOT a command/arg)
btn = None
if not text.startswith("/"):
if text.startswith("➕"):
btn = "add"
elif text.startswith("\U0001F5D1"):
btn = "del"
elif text.startswith("\U0001F4CB"):
btn = "list"
elif text.startswith("\U0001F4CA"):
btn = "status"
elif text.startswith("❓"):
btn = "help"
# conversation: awaiting a client name to create
if _bot_state.get(chat) == "add":
if btn is None and not text.startswith("/"):
_bot_state.pop(chat, None)
bot_send_new_client(chat, text)
continue
_bot_state.pop(chat, None) # a button/command cancels the prompt
# parse slash command + arg (buttons carry no arg)
cmd = ""
arg = ""
if text.startswith("/"):
toks = text.split()
cmd = toks[0].split("@")[0].lower()
arg = toks[1] if len(toks) > 1 else ""
if btn == "add" or cmd in ("/add", "/new", "/create"):
if arg:
bot_send_new_client(chat, arg)
else:
_bot_state[chat] = "add"
tg_api("sendMessage", {"chat_id": chat, "text": "Введите имя нового клиента (латиница, цифры, _ и -):"})
elif btn == "del" or cmd in ("/del", "/delete", "/remove", "/rm"):
if arg:
_, rep = delete_client_cli(arg)
tg_api("sendMessage", {"chat_id": chat, "text": rep, "parse_mode": "HTML", "reply_markup": MAIN_KB})
else:
cl = get_clients()
if not cl:
tg_api("sendMessage", {"chat_id": chat, "text": "Клиентов нет.", "reply_markup": MAIN_KB})
else:
kb = {"inline_keyboard": [[{"text": "\U0001F5D1 " + c.get("client_id", ""), "callback_data": "del:" + c.get("client_id", "")}] for c in cl]}
tg_api("sendMessage", {"chat_id": chat, "text": "Кого удалить? Нажми на клиента:", "reply_markup": kb})
elif btn == "list" or cmd in ("/list", "/ls"):
cl = [c.get("client_id", "") for c in get_clients()]
rep = ("Клиенты (%d):\n" % len(cl)) + "\n".join("• " + c for c in cl) if cl else "Клиентов нет."
tg_api("sendMessage", {"chat_id": chat, "text": rep, "reply_markup": MAIN_KB})
elif btn == "status" or cmd in ("/status", "/stat"):
tg_api("sendMessage", {"chat_id": chat, "text": _bot_status_text(), "parse_mode": "HTML", "reply_markup": MAIN_KB})
elif btn == "help" or cmd in ("/help", "/start"):
rep = ("PCA Phobos — бот\n\n"
"Кнопки внизу или команды:\n"
"➕ /add <имя> — создать клиента\n"
"\U0001F5D1 /del <имя> — удалить\n"
"\U0001F4CB /list — список\n"
"\U0001F4CA /status — статус")
tg_api("sendMessage", {"chat_id": chat, "text": rep, "parse_mode": "HTML", "reply_markup": MAIN_KB})
except Exception:
time.sleep(10)
def get_wg_peers():
"""Parse `wg show wg0` to get active peers with transfer/handshake info."""
try:
out = subprocess.check_output(["wg", "show", "wg0"], text=True, timeout=5)
except Exception:
return {}
peers = {}
current_pub = None
for line in out.split("\n"):
line = line.strip()
if line.startswith("peer:"):
current_pub = line.split("peer:")[1].strip()
peers[current_pub] = {}
elif current_pub and ":" in line:
key, val = line.split(":", 1)
peers[current_pub][key.strip()] = val.strip()
return peers
def get_clients():
"""Read all clients from /opt/Phobos/clients/*/metadata.json."""
clients = []
clients_path = Path(CLIENTS_DIR)
if not clients_path.exists():
return clients
for d in sorted(clients_path.iterdir()):
meta_file = d / "metadata.json"
if meta_file.exists():
try:
with open(meta_file) as f:
meta = json.load(f)
meta["_dir"] = str(d)
clients.append(meta)
except Exception:
pass
return clients
def get_active_sessions():
"""Combine WG peers with client metadata to build session list."""
peers = get_wg_peers()
clients = get_clients()
pub_to_client = {}
for c in clients:
pub_to_client[c.get("public_key", "")] = c
sessions = []
for pub_key, info in peers.items():
handshake = info.get("latest handshake", "")
if not handshake:
continue
client = pub_to_client.get(pub_key, {})
client_id = client.get("client_id", "unknown")
tunnel_ip = client.get("tunnel_ip_v4", "")
endpoint = info.get("endpoint", "")
real_ip = endpoint.split(":")[0] if endpoint else ""
rx = info.get("transfer", "")
rx_bytes = rx.split("received,")[0].strip() if "received," in rx else ""
tx_bytes = rx.split("received,")[1].strip().replace("sent", "").strip() if "received," in rx else ""
sessions.append({
"client_id": client_id,
"public_key": pub_key,
"tunnel_ip": tunnel_ip,
"real_ip": real_ip,
"endpoint": endpoint,
"handshake": handshake,
"rx": rx_bytes,
"tx": tx_bytes,
})
return sessions
def is_peer_online(handshake_str):
"""Check if peer had a handshake within last 3 minutes."""
try:
parts = handshake_str.split(",")
total_seconds = 0
for p in parts:
p = p.strip()
if "minute" in p:
total_seconds += int(re.search(r"(\d+)", p).group(1)) * 60
elif "second" in p:
total_seconds += int(re.search(r"(\d+)", p).group(1))
elif "hour" in p:
total_seconds += int(re.search(r"(\d+)", p).group(1)) * 3600
return total_seconds < 180
except Exception:
return False
def kick_peer(public_key):
"""Remove and re-add peer to force disconnect."""
try:
out = subprocess.check_output(["wg", "show", "wg0"], text=True, timeout=5)
allowed = ""
found = False
for line in out.split("\n"):
if line.strip().startswith("peer:") and public_key in line:
found = True
elif found and "allowed ips:" in line:
allowed = line.split("allowed ips:")[1].strip()
break
subprocess.run(["wg", "set", "wg0", "peer", public_key, "remove"], timeout=5)
if allowed:
subprocess.run(["wg", "set", "wg0", "peer", public_key, "allowed-ips", allowed], timeout=5)
return True
except Exception:
return False
def check_expiry():
"""Check subscription expiry, lock expired clients."""
s = load_settings() # Always reload fresh to avoid overwriting concurrent changes
subs = s.get("subscriptions", {})
today = datetime.now().date()
changed = False
for client_id, info in list(subs.items()):
if not info.get("expiry"):
continue
try:
exp_date = datetime.strptime(info["expiry"], "%Y-%m-%d").date()
except ValueError:
continue
days_left = (exp_date - today).days
if days_left <= 0 and not info.get("locked"):
info["locked"] = True
changed = True
kick_client_by_id(client_id)
tg_send(f"⛔ {client_id} — подписка истекла! Клиент заблокирован.")
elif days_left == 3 and not info.get("warn3"):
info["warn3"] = True
changed = True
tg_send(f"⚠️ {client_id} — подписка истекает через 3 дня ({info['expiry']})")
elif days_left == 1 and not info.get("warn1"):
info["warn1"] = True
changed = True
tg_send(f"⚠️ {client_id} — подписка истекает ЗАВТРА ({info['expiry']})")
if changed:
save_settings(s)
def kick_client_by_id(client_id):
"""Find client's public key and kick them."""
clients = get_clients()
for c in clients:
if c.get("client_id") == client_id:
kick_peer(c.get("public_key", ""))
return True
return False
prev_session_keys = None
prev_server_status = {}
_down_streak = {}
_last_fanout = 0
server_stats_cache = {}
server_handshakes_cache = {}
def count_client_peers():
"""Count only client peers (exclude secondary server peers)."""
clients = get_clients()
client_pubs = {c.get("public_key", "") for c in clients}
wg_peers = get_wg_peers()
return sum(1 for pub in wg_peers if pub in client_pubs)
def get_local_stats():
try:
cpu = subprocess.getoutput("top -bn1 | grep 'Cpu(s)' | awk '{print $2}'").strip()
mem = subprocess.getoutput("free -m | awk '/Mem:/{printf \"%.0f/%dMB\", $3, $2}'").strip()
peers = count_client_peers()
return {"cpu": cpu + "%", "mem": mem, "peers": peers, "status": "ok"}
except Exception:
return {"cpu": "?", "mem": "?", "peers": 0, "status": "ok"}
def get_remote_stats(server):
import urllib.request
url = f"http://{server['ip']}:8444/api/health"
req = urllib.request.Request(url, headers={"X-API-Key": server.get("api_key", "")})
data = None
for _attempt in range(2): # retry once: smooths transient timeouts (1-worker agent)
try:
data = json.loads(urllib.request.urlopen(req, timeout=6).read())
break
except Exception:
data = None
try:
if data is None:
raise Exception("no response")
# cache handshakes so page renders never block on remote HTTP
server_handshakes_cache[server.get("ip", "")] = data.get("handshakes", {}) or {}
# Filter peers: only count known client public keys
clients = get_clients()
client_pubs = {c.get("public_key", "") for c in clients}
remote_keys = data.get("peer_keys", [])
client_peers = sum(1 for k in remote_keys if k in client_pubs)
return {"cpu": data.get("cpu", "—"), "mem": data.get("mem", "—"), "peers": client_peers, "status": "ok" if data.get("status") == "ok" else "down"}
except Exception:
return {"cpu": "—", "mem": "—", "peers": 0, "status": "down"}
def server_load(stats):
"""Estimate server load 0..1 from cached stats (CPU+RAM). Missing stats =
neutral 0.5; explicitly down = 1.0 (avoid). Used by load-aware rebalance."""
if not stats:
return 0.5
if stats.get("status") == "down":
return 1.0
try:
cpu = float(str(stats.get("cpu", "")).replace("%", "").strip()) / 100.0
except Exception:
cpu = 0.5
mem = 0.5
try:
used, total = str(stats.get("mem", "")).replace("MB", "").split("/")
mem = float(used) / max(float(total), 1.0)
except Exception:
pass
return max(0.0, min(1.0, 0.6 * cpu + 0.4 * mem))
def remote_handshakes(server_ip):
# Pure cache read — the background session_monitor refreshes it every cycle.
# Page renders must NEVER block on remote HTTP (that caused multi-second hangs).
return server_handshakes_cache.get(server_ip, {})
def auto_assign_server(client_id):
s = load_settings()
assignments = s.get("client_assignments", {})
if client_id in assignments:
return assignments[client_id]
all_servers = get_all_servers_ordered()
server_ips = [sv["ip"] for sv in all_servers]
if not server_ips:
return SERVER_IP
counts = {ip: 0 for ip in server_ips}
for cid, sip in assignments.items():
if sip in counts:
counts[sip] += 1
assigned = min(counts, key=counts.get)
assignments[client_id] = assigned
s["client_assignments"] = assignments
save_settings(s)
return assigned
def add_peer_to_server(server_ip, public_key, allowed_ips):
"""Add WG peer to a server via its API (secondary) or locally (primary)."""
if server_ip == SERVER_IP:
# Primary — add locally
try:
subprocess.run(["wg", "set", "wg0", "peer", public_key, "allowed-ips", allowed_ips], check=True, timeout=5)
subprocess.run("wg-quick save wg0", shell=True, timeout=5)
return {"status": "ok"}
except Exception as e:
return {"status": "error", "msg": str(e)}
# Secondary — call API
servers = load_servers()
api_key = ""
for srv in servers:
if srv["ip"] == server_ip:
api_key = srv.get("api_key", "")
break
if not api_key:
return {"status": "error", "msg": f"No API key for {server_ip}"}
try:
import urllib.request
data = json.dumps({"public_key": public_key, "allowed_ips": allowed_ips}).encode()
req = urllib.request.Request(
f"http://{server_ip}:8444/api/peers/add",
data=data,
headers={"X-API-Key": api_key, "Content-Type": "application/json"},
method="POST"
)
resp = urllib.request.urlopen(req, timeout=10)
return json.loads(resp.read())
except Exception as e:
return {"status": "error", "msg": str(e)[:200]}
def restart_router_obfuscator(client_id):
"""Restart wg-obfuscator on router via SSH (tries tunnel, then static)."""
s = load_settings()
access = s.get("router_access", {}).get(client_id)
if not access:
return {"status": "error", "msg": f"No SSH for {client_id}"}
ok_ip, out = _ssh_run(access, "/opt/etc/init.d/S49wg-obfuscator restart 2>/dev/null; echo OK")
if ok_ip and "OK" in out:
return {"status": "ok"}
return {"status": "error", "msg": out[:200]}
def generate_failover_conf_for_client(client_id):
s = load_settings()
assignments = s.get("client_assignments", {})
assigned_ip = assignments.get(client_id, SERVER_IP)
all_servers = get_all_servers_ordered()
by_ip = {sv["ip"]: sv for sv in all_servers}
ordered = []
if assigned_ip in by_ip:
ordered.append(by_ip[assigned_ip])
for sv in all_servers:
if sv["ip"] != assigned_ip:
ordered.append(sv)
# Get main server WG public key
try:
main_wg_pub = subprocess.check_output(["wg", "show", "wg0", "public-key"], text=True, timeout=5).strip()
except Exception:
main_wg_pub = ""
lines = ["# Phobos Failover Configuration"]
for i, srv in enumerate(ordered, 1):
lines.append(f"SERVER_{i}={srv['ip']}:{srv.get('ports', '2083,5443,993')}")
lines.append(f"KEY_{i}={srv.get('obfuscator_key', '')}")
wg_pub = main_wg_pub if srv.get("is_primary") else srv.get("wg_public_key", "")
if wg_pub:
lines.append(f"WGKEY_{i}={wg_pub}")
return "\n".join(lines) + "\n"
def panel_version():
try:
return open(os.path.join(PANEL_DIR, ".version")).read().strip()
except Exception:
return "unknown"
def run_update(arg):
import subprocess
try:
r = subprocess.run(["/opt/Phobos/server/update.sh", arg],
capture_output=True, text=True, timeout=150)
return (r.stdout + r.stderr)[-1200:]
except Exception as e:
return "update error: " + str(e)
def build_phone_config(cid, mode="android"):
"""Assemble a phone client config from the client's files.
android = WireGuard + [instance] obfuscator -> phobos:// link (PhobosWG app).
ios = plain WireGuard (no obfuscation), Endpoint -> server:51820."""
import base64, urllib.parse
cdir = os.path.join(CLIENTS_DIR, cid)
wgpath = os.path.join(cdir, cid + ".conf")
instpath = os.path.join(cdir, "wg-obfuscator.conf")
if not os.path.exists(wgpath):
return None, None
wg = open(wgpath).read().rstrip()
if mode == "ios":
server_ip = SERVER_IP
try:
for ln in open(instpath):
if ln.strip().startswith("target"):
server_ip = ln.split("=", 1)[1].strip().split(":")[0]
break
except Exception:
pass
out = []
for ln in wg.split("\n"):
out.append("Endpoint = %s:51820" % server_ip if ln.strip().startswith("Endpoint") else ln)
return "\n".join(out) + "\n", None
inst = ""
if os.path.exists(instpath):
inst = open(instpath).read().strip()
conf = wg + "\n\n" + inst + "\n"
b64 = base64.urlsafe_b64encode(conf.encode()).decode().rstrip("=")
link = "phobos://" + b64 + "#" + urllib.parse.quote(cid)
return conf, link
def qr_datauri(data):
try:
import qrcode, io, base64
buf = io.BytesIO()
qrcode.make(data).save(buf, format="PNG")
return "data:image/png;base64," + base64.b64encode(buf.getvalue()).decode()
except Exception:
return ""
def fanout_router_config(client_id):
"""Push this client's failover.conf to every secondary server's agent so
routers can PULL it through the tunnel (10.25.0.1:8444) — survives a public
panel-IP ban and follows the active server. Best-effort, non-blocking."""
try:
conf = generate_failover_conf_for_client(client_id)
except Exception:
return
import urllib.request
payload = json.dumps({"client_id": client_id, "conf": conf}).encode()
for srv in load_servers():
ip = srv.get("ip", "")
key = srv.get("api_key", "")
if not ip:
continue
# skip servers the monitor already knows are down (avoid blocking)
if server_stats_cache.get(ip, {}).get("status") == "down":
continue
try:
req = urllib.request.Request(
f"http://{ip}:8444/api/router-config-set",
data=payload,
headers={"X-API-Key": key, "Content-Type": "application/json"},
method="POST")
urllib.request.urlopen(req, timeout=3)
except Exception:
pass
def _ssh_run(access, cmd, timeout=20):
"""Try SSH: tunnel IP first, then static IP as fallback."""
ssh_user = access.get("ssh_user", "root")
ssh_pass = access.get("ssh_pass", "")
ips_to_try = []
if access.get("ssh_ip"):
ips_to_try.append(access["ssh_ip"])
if access.get("ssh_static"):
ips_to_try.append(access["ssh_static"])
if not ips_to_try or not ssh_pass:
return None, "SSH IP/pass missing"
for ip in ips_to_try:
try:
result = subprocess.run(
["sshpass", "-p", ssh_pass, "ssh",
"-o", "StrictHostKeyChecking=no", "-o", "ConnectTimeout=3",
f"{ssh_user}@{ip}", cmd],
capture_output=True, text=True, timeout=timeout
)
if result.returncode == 0:
return ip, result.stdout
except Exception:
continue
return None, f"SSH failed on all IPs: {', '.join(ips_to_try)}"
def push_config_to_router(client_id):
s = load_settings()
access = s.get("router_access", {}).get(client_id)
if not access:
return {"status": "error", "msg": f"No SSH for {client_id}"}
conf = generate_failover_conf_for_client(client_id)
write_cmd = f"mkdir -p /opt/etc/Phobos && cat > /opt/etc/Phobos/failover.conf << 'FAILCONF'\n{conf}FAILCONF"
ok_ip, out = _ssh_run(access, write_cmd)
if ok_ip:
return {"status": "ok", "msg": f"Config pushed to {client_id} ({ok_ip})"}
return {"status": "error", "msg": out[:200]}
def session_monitor():
global prev_session_keys, prev_server_status, server_stats_cache, _last_fanout, _down_streak
while True:
try:
s = load_settings()
interval = s.get("monitor_interval", 30)
# Client session monitoring
sessions = get_active_sessions()
current_keys = set()
for sess in sessions:
if is_peer_online(sess.get("handshake", "")):
key = (sess["client_id"], sess["real_ip"])
current_keys.add(key)
if prev_session_keys is not None:
labels = s.get("labels", {})
for key in current_keys - prev_session_keys:
client_id, real_ip = key
label = labels.get(real_ip, "")
name = f"{label} ({real_ip})" if label else real_ip
if notify_on(client_id):
tg_send(f"🟢 {client_id} connected — {name}")
for key in prev_session_keys - current_keys:
client_id, real_ip = key
label = labels.get(real_ip, "")
name = f"{label} ({real_ip})" if label else real_ip
if notify_on(client_id):
tg_send(f"🔴 {client_id} disconnected — {name}")
prev_session_keys = current_keys
# Server health monitoring
servers = load_servers()
local_stats = get_local_stats()
server_stats_cache[SERVER_IP] = local_stats
DOWN_CONFIRM = 3 # consecutive failed cycles before declaring a server DOWN
for srv in servers:
ip = srv["ip"]
stats = get_remote_stats(srv)
raw = stats.get("status", "down")
streak = _down_streak.get(ip, 0)
streak = streak + 1 if raw == "down" else 0
_down_streak[ip] = streak
# update cache on OK, or only once a DOWN is confirmed (keep last-good on transient)
if raw == "ok" or streak >= DOWN_CONFIRM:
server_stats_cache[ip] = stats
eff = "down" if streak >= DOWN_CONFIRM else "ok"
was_up = prev_server_status.get(ip, "ok")
if was_up == "ok" and eff == "down":
tg_send(f"🔴 Server {ip} DOWN!")
elif was_up == "down" and eff == "ok":
tg_send(f"🟢 Server {ip} back ONLINE")
prev_server_status[ip] = eff
# periodic config fan-out to secondaries (covers CLI-created clients
# and config drift; gated ~5 min, skips down servers)
global _last_fanout
if time.time() - _last_fanout >= 300:
_last_fanout = time.time()
for _c in get_clients():
fanout_router_config(_c.get("client_id", ""))
check_expiry()
time.sleep(interval)
except Exception:
time.sleep(30)
monitor_thread = threading.Thread(target=session_monitor, daemon=True)
monitor_thread.start()
bot_thread = threading.Thread(target=telegram_bot, daemon=True)
bot_thread.start()
PAGE = """
")
def set_lang(code):
from flask import make_response
resp = make_response(redirect(request.referrer or url_for("dashboard")))
resp.set_cookie("lang", "en" if code == "en" else "ru", max_age=31536000)
return resp
@app.route("/login", methods=["GET", "POST"])
def login():
s = load_settings()
msg = ""
if request.method == "POST":
if request.form.get("password") == s["admin_pass"]:
session["auth"] = True
return redirect(url_for("dashboard"))
msg = '{tr("Неверный пароль","Wrong password")}'
html = f"""
{tr("Вход в панель","Sign in")}
{msg}
"""
return render(html)
@app.route("/logout")
def logout():
session.clear()
return redirect(url_for("login"))
def auth_required(f):
from functools import wraps
@wraps(f)
def decorated(*args, **kwargs):
if not session.get("auth"):
return redirect(url_for("login"))
return f(*args, **kwargs)
return decorated
@app.route("/")
@auth_required
def dashboard():
return redirect(url_for("sessions_page"))
@app.route("/sessions")
@auth_required
def sessions_page():
s = load_settings()
labels = s.get("labels", {})
sessions = get_active_sessions()
# Online = fresh handshake on ANY server (router may be on a backup)
online_all = set()
for _pub, _info in get_wg_peers().items():
if is_peer_online(_info.get("latest handshake", "")):
online_all.add(_pub)
for _ip, _hs in server_handshakes_cache.items():
for _pub, _ts in _hs.items():
if _ts and (time.time() - _ts) < 180:
online_all.add(_pub)
rows = ""
online_count = 0
for sess in sessions:
online = sess["public_key"] in online_all
if online:
online_count += 1
status = 'Online' if online else 'Offline'
real_ip = sess["real_ip"]
label = labels.get(sess.get("tunnel_ip", ""), "") or labels.get(real_ip, "")
rows += f"""
{sess['client_id']}
{sess['tunnel_ip']}
{real_ip}
{('' + label + '') if label else '—'}
{sess['handshake']}
{sess['rx']} / {sess['tx']}
{status}
{tr("Отключить","Kick")}{hlp("Принудительно разорвать текущую сессию клиента (сбросить WG-пир). Клиент переподключится автоматически.","Force-drop the client current session (reset the WG peer). The client reconnects automatically.")}
"""
if not rows:
rows = '{tr("Нет активных сессий","No active sessions")} '
html = f"""
{nav('sessions', online_count)}
{tr("Активные сессии","Active sessions")}
Клиент VPN IP Real IP Метка Handshake RX / TX Статус
{rows}
"""
return render(html)
@app.route("/kick/")
@auth_required
def kick(pub_key):
kick_peer(pub_key)
return redirect(url_for("sessions_page"))
@app.route("/clients", methods=["GET", "POST"])
@auth_required
def clients_page():
s = load_settings()
subs = s.get("subscriptions", {})
msg = ""
if request.method == "POST":
action = request.form.get("action")
if action == "add":
name = request.form.get("name", "").strip()
if name and re.match(r"^[a-zA-Z0-9_-]+$", name):
try:
out = subprocess.check_output(
[f"{PHOBOS_DIR}/repo/server/scripts/phobos-client.sh", "add", name],
text=True, timeout=30, stderr=subprocess.STDOUT,
env={**os.environ, **_load_server_env()}
)
auto_assign_server(name)
fanout_router_config(name) # push conf to all servers (tunnel pull)
msg = f'Client {name} created'
except subprocess.CalledProcessError as e:
msg = f'{e.output}'
else:
msg = 'Имя: буквы, цифры, _ и -'
elif action == "delete":
client_id = request.form.get("client_id", "").strip()
if client_id:
try:
out = subprocess.check_output(
[f"{PHOBOS_DIR}/repo/server/scripts/phobos-client.sh", "remove", client_id],
text=True, timeout=15, stderr=subprocess.STDOUT,
env={**os.environ, **_load_server_env()}
)
msg = f'Клиент {client_id} удалён'
except subprocess.CalledProcessError as e:
msg = f'{e.output}'
elif action == "toggle_notify":
client_id = request.form.get("client_id", "").strip()
if client_id:
cn = s.setdefault("client_notify", {})
cn[client_id] = not cn.get(client_id, True)
save_settings(s)
st = "вкл" if cn[client_id] else "выкл"
msg = f'Уведомления {client_id}: {st}'
elif action == "set_expiry":
client_id = request.form.get("client_id", "").strip()
expiry = request.form.get("expiry", "").strip()
if client_id:
if client_id not in subs:
subs[client_id] = {}
subs[client_id]["expiry"] = expiry
subs[client_id].pop("locked", None)
subs[client_id].pop("warn3", None)
subs[client_id].pop("warn1", None)
save_settings(s)
msg = f'Срок для {client_id} обновлён'
elif action == "push_config":
client_id = request.form.get("client_id", "").strip()
res = push_config_to_router(client_id)
if res["status"] == "ok":
msg = f'{res["msg"]}'
else:
msg = f'{client_id}: saved — router will pull config within ~2 min (SSH not needed for NAT routers)'
elif action == "assign_server":
client_id = request.form.get("client_id", "").strip()
server_ip = request.form.get("server_ip", "").strip()
if client_id and server_ip:
# 1. Find client public key + tunnel IP for WG peer
client_meta = None
for c in get_clients():
if c.get("client_id") == client_id:
client_meta = c
break
errors = []
# 2. Add WG peer on target server
if client_meta:
pub = client_meta.get("public_key", "")
tip = client_meta.get("tunnel_ip_v4", "").split("/")[0]
if pub and tip:
allowed = f"{tip}/32"
peer_res = add_peer_to_server(server_ip, pub, allowed)
if peer_res.get("status") != "ok":
errors.append(f"Peer add: {peer_res.get('msg', 'fail')}")
# 3. Save assignment
ca = s.get("client_assignments", {})
ca[client_id] = server_ip
s["client_assignments"] = ca
save_settings(s)
fanout_router_config(client_id) # push conf to all servers (tunnel pull)
# 4. Router applies via pull (~15s). No synchronous SSH push:
# NAT routers can't be reached, and the SSH attempt blocked
# the click for seconds. Pull channel handles the apply.
if errors:
msg = f'{client_id} → {server_ip}: {"; ".join(errors)}'
else:
msg = f'{client_id} → {server_ip} ✓ (saved — router applies within ~15s)'
clients = get_clients()
peers = get_wg_peers()
online_pubs = set()
for pub, info in peers.items():
if is_peer_online(info.get("latest handshake", "")):
online_pubs.add(pub)
all_srv = get_all_servers_ordered()
assignments = s.get("client_assignments", {})
assignments_notify = s.get("client_notify", {})
# Online = fresh handshake on ANY server (truthful during failover/transition)
online_all = set(online_pubs)
for _sv in all_srv:
if _sv.get("ip") == SERVER_IP:
continue
for _pub, _ts in remote_handshakes(_sv["ip"]).items():
if _ts and (time.time() - _ts) < 180:
online_all.add(_pub)
rows = ""
today = datetime.now().date()
for c in clients:
cid = c.get("client_id", "")
pub = c.get("public_key", "")
ip = c.get("tunnel_ip_v4", "")
created = c.get("created_at", "")[:10]
is_online = pub in online_all
status = 'Online' if is_online else 'Offline'
sub = subs.get(cid, {})
expiry = sub.get("expiry", "")
locked = sub.get("locked", False)
expiry_badge = ""
if expiry:
try:
exp_date = datetime.strptime(expiry, "%Y-%m-%d").date()
days = (exp_date - today).days
if locked:
expiry_badge = 'expired'
elif days <= 1:
expiry_badge = f'{days}d'
elif days <= 3:
expiry_badge = f'{days}d'
else:
expiry_badge = f'{days}d'
except ValueError:
pass
assigned_ip = assignments.get(cid, SERVER_IP)
srv_opts = ""
for sv in all_srv:
sel = " selected" if sv["ip"] == assigned_ip else ""
label = f'{sv["ip"]} (Primary)' if sv.get("is_primary") else sv["ip"]
srv_opts += f''
set_lbl = tr("Задать", "Set")
h_set = hlp("Выбери сервер выхода для этого роутера и нажми «Задать». Роутер сам переключится за ~15 секунд. Кнопку Push нажимать НЕ нужно — конфигурация подтягивается автоматически.",
"Pick the exit server for this router and press Set. The router switches itself within ~15 seconds. You do NOT need a Push button — the config is pulled automatically.")
h_exp = hlp("Дата окончания подписки клиента. После этой даты клиент автоматически блокируется. Оставь пустым — без срока.",
"Client subscription expiry date. After this date the client is locked automatically. Leave empty for no expiry.")
h_del = hlp("Безвозвратно удалить клиента, его ключи и конфигурацию. Отменить нельзя.",
"Permanently delete the client, its keys and config. Cannot be undone.")
del_confirm = tr("Удалить " + cid + "?", "Delete " + cid + "?")
_notify = assignments_notify.get(cid, True)
notify_icon = "🔔" if _notify else "🔕"
notify_bg = "#4f46e5" if _notify else "#64748b"
notify_ttl = tr("Уведомления в Telegram включены — нажми, чтобы выключить", "Telegram notifications ON — click to turn off") if _notify else tr("Уведомления выключены — нажми, чтобы включить", "Notifications OFF — click to turn on")
rows += f"""
{cid}
{ip}
{status}
📱
🍎
{h_del}
"""
html = f"""
{nav('clients')}
{msg}
{tr("VPN-клиенты", "VPN Clients")}
{tr("Клиент","Client")} VPN IP {tr("Статус","Status")} {tr("Сервер","Server")} {tr("Срок","Expiry")} {tr("Действия","Actions")}
{rows}
"""
return render(html)
@app.route("/labels", methods=["GET", "POST"])
@auth_required
def labels_page():
s = load_settings()
msg = ""
if request.method == "POST":
action = request.form.get("action")
if action == "add":
ip = request.form.get("ip", "").strip()
label = request.form.get("label", "").strip()
if ip and label:
s["labels"][ip] = label
save_settings(s)
msg = f'Метка добавлена: {ip} → {label}'
elif action == "delete":
ip = request.form.get("ip", "").strip()
s["labels"].pop(ip, None)
save_settings(s)
msg = f'Метка удалена'
labels = s.get("labels", {})
rows = ""
for ip, label in sorted(labels.items()):
rows += f"""
{ip} {label}
"""
html = f"""
{nav('labels')}
{msg}
{tr("Метки (по IP)","Labels (by IP)")} {hlp("Человекочитаемые имена для IP/туннелей — показываются в Сессиях вместо голого адреса.","Human-readable names for IPs/tunnels — shown in Sessions instead of the raw address.")}
Привяжите понятное имя (квартира, офис, дача) к внешнему IP адресу роутера. Метка отображается в таблице сессий и в Telegram уведомлениях вместо голого IP.
Real IP Метка
{rows}
"""
return render(html)
@app.route("/settings", methods=["GET", "POST"])
@auth_required
def settings_page():
s = load_settings()
msg = ""
if request.method == "POST" and request.form.get("action") in ("update", "rollback", "check"):
_act = request.form.get("action")
if _act == "update":
_out = run_update(request.form.get("target", "").strip() or "main")
elif _act == "check":
_out = run_update("--check")
else:
_out = run_update("--rollback")
s = load_settings()
return render(f'''{nav("settings")}''')
if request.method == "POST":
new_pass = request.form.get("admin_pass", "").strip()
if new_pass:
s["admin_pass"] = new_pass
s["tg_bot_token"] = request.form.get("tg_bot_token", "").strip()
s["tg_chat_id"] = request.form.get("tg_chat_id", "").strip()
try:
s["monitor_interval"] = max(10, int(request.form.get("monitor_interval", 30)))
except ValueError:
pass
save_settings(s)
msg = 'Настройки сохранены'
html = f"""
{nav('settings')}
{msg}
{tr("Настройки панели","Panel settings")}
{tr("Версия и обновления","Version & updates")} {hlp("Скачивает свежие файлы PCA-слоя (панель + скрипты) с GitHub и перезапускает панель. Ключи WireGuard, server.env и клиенты НЕ трогаются. Пусто или main = последняя версия; можно указать тег (например v1.1.0). Откат восстанавливает предыдущую версию из автоматического бэкапа.", "Fetches the latest PCA files (panel + scripts) from GitHub and restarts the panel. WireGuard keys, server.env and clients are NOT touched. Empty or main = latest; you can pass a tag (e.g. v1.1.0). Rollback restores the previous version from an automatic backup.")}
{tr("Установленная версия","Installed version")}: {panel_version()}
{tr("Информация о сервере","Server info")}
VPS IP {SERVER_IP}
{tr("WireGuard порт","WireGuard port")} 51820 (localhost)
{tr("Порты обфускатора","Obfuscator ports")} {_load_server_env().get('OBFUSCATOR_PORTS', '2083,5443,993')}
{tr("Порт панели","Panel port")} 8443
{tr("Клиенты Phobos","Phobos clients")} {CLIENTS_DIR}
{tr("Установка на роутер","Router install")}
{tr("Keenetic/Netcraze с Entware — выполнить по SSH на роутере:","Keenetic/Netcraze with Entware — run over SSH on the router:")}
{_render_install_commands()}
"""
return render(html)
TOKENS_FILE = f"{PHOBOS_DIR}/tokens/tokens.json"
def _render_install_commands():
"""Read tokens and render install commands for each client."""
try:
if os.path.exists(TOKENS_FILE):
with open(TOKENS_FILE) as f:
tokens = json.load(f)
else:
tokens = []
except Exception:
tokens = []
if not tokens:
return 'Нет активных токенов. Добавьте клиента.
'
lines = ""
for t in tokens:
client = t.get("client", "?")
token = t.get("token", "")
expires = t.get("expires", 0)
exp_str = datetime.fromtimestamp(expires).strftime("%Y-%m-%d %H:%M") if expires else "?"
cmd = f"wget -O - http://{SERVER_IP}/init/{token}.sh | sh"
lines += f"""
{client}
(до {exp_str})
{cmd}
"""
return lines
SERVERS_FILE = f"{PANEL_DIR}/servers.json"
def load_servers():
if os.path.exists(SERVERS_FILE):
with open(SERVERS_FILE) as f:
return json.load(f)
return []
def save_servers(servers):
with open(SERVERS_FILE, "w") as f:
json.dump(servers, f, indent=2)
def get_main_server_info():
env = _load_server_env()
return {
"ip": SERVER_IP,
"ports": env.get("OBFUSCATOR_PORTS", "2083,5443,993"),
"obfuscator_key": env.get("OBFUSCATOR_KEY", ""),
"is_primary": True
}
def get_all_servers_ordered():
s = load_settings()
servers = load_servers()
main_info = get_main_server_info()
order = s.get("server_order", [])
by_ip = {}
by_ip[main_info["ip"]] = {"ip": main_info["ip"], "ports": main_info["ports"],
"obfuscator_key": main_info["obfuscator_key"], "is_primary": True}
for srv in servers:
by_ip[srv["ip"]] = {**srv, "is_primary": False}
if not order or main_info["ip"] not in order:
order = [main_info["ip"]] + [sv["ip"] for sv in servers]
result = []
for ip in order:
if ip in by_ip:
result.append(by_ip[ip])
for ip, sv in by_ip.items():
if ip not in order:
result.append(sv)
return result
def check_server_health(server):
"""Check secondary server health via API."""
try:
import urllib.request
url = f"http://{server['ip']}:8444/api/health"
req = urllib.request.Request(url, headers={"X-API-Key": server.get("api_key", "")})
resp = urllib.request.urlopen(req, timeout=5)
return json.loads(resp.read())
except Exception as e:
return {"status": "error", "error": str(e)}
def sync_peer_to_server(server, public_key, allowed_ips, action="add"):
"""Add or remove peer on secondary server."""
try:
import urllib.request
url = f"http://{server['ip']}:8444/api/peers/{action}"
data = json.dumps({"public_key": public_key, "allowed_ips": allowed_ips}).encode()
req = urllib.request.Request(url, data=data, headers={
"Content-Type": "application/json",
"X-API-Key": server.get("api_key", "")
})
resp = urllib.request.urlopen(req, timeout=10)
return json.loads(resp.read())
except Exception:
return {"status": "error"}
def sync_peer_to_all_servers(public_key, allowed_ips, action="add"):
"""Sync peer to all secondary servers."""
servers = load_servers()
for srv in servers:
if srv.get("enabled", True):
sync_peer_to_server(srv, public_key, allowed_ips, action)
@app.route("/client//phone/")
@auth_required
def client_phone(cid, mode):
if mode not in ("android", "ios"):
mode = "android"
conf, link = build_phone_config(cid, mode)
if conf is None:
return ("client not found", 404)
qr = qr_datauri(link if (mode == "android" and link) else conf)
if mode == "android":
title = tr("Android — PhobosWG (с обфускацией)", "Android — PhobosWG (obfuscated)")
apk_btn = ""
if os.path.exists("/opt/Phobos/www/app/PhobosWG.apk"):
apk_btn = (f'📥 {tr("Скачать приложение PhobosWG (APK)","Download PhobosWG app (APK)")}
')
extra = (apk_btn +
f'{tr("phobos:// ссылка — импорт в PhobosWG:", "phobos:// link — import into PhobosWG:")}
'
f'')
else:
title = tr("iPhone / iOS WireGuard (без обфускации)", "iPhone / iOS WireGuard (plain, no obfuscation)")
extra = (f'📥 {tr("WireGuard в App Store","WireGuard on the App Store")}
'
f'{tr("⚠ Обычный WireGuard без маскировки (Endpoint → :51820). На сервере должен быть открыт порт 51820 (ALLOW_PLAIN_WG=1).", "⚠ Plain WireGuard, no obfuscation (Endpoint → :51820). Port 51820 must be open on the server (ALLOW_PLAIN_WG=1).")}
')
html = f"""
{nav('clients')}
{cid} — {title}

{tr("Отсканируй QR в приложении или импортируй конфиг:", "Scan the QR in the app or import the config:")}
{extra}
"""
return render(html)
@app.route("/api/router-config/")
def router_config(client_id):
"""NAT-friendly config pull: router fetches its own failover.conf.
Auth via per-client pull_token (falls back to server_api_key)."""
st = load_settings()
token = request.args.get("token", "")
acc = st.get("router_access", {}).get(client_id, {})
expected = acc.get("pull_token", "") or st.get("server_api_key", "")
if not expected or token != expected:
return ("forbidden", 403)
conf = generate_failover_conf_for_client(client_id)
return (conf, 200, {"Content-Type": "text/plain; charset=utf-8"})
@app.route("/api/servers/register", methods=["POST"])
def api_register_server():
s = load_settings()
api_key = request.headers.get("X-API-Key", "")
if api_key != s.get("server_api_key", ""):
return json.dumps({"error": "unauthorized"}), 401, {"Content-Type": "application/json"}
data = request.json
servers = load_servers()
ip = data.get("ip", "")
for srv in servers:
if srv["ip"] == ip:
srv.update(data)
srv["last_seen"] = datetime.now().isoformat()
save_servers(servers)
return json.dumps({"status": "updated"}), 200, {"Content-Type": "application/json"}
data["last_seen"] = datetime.now().isoformat()
data["enabled"] = True
data["api_key"] = api_key
servers.append(data)
save_servers(servers)
return json.dumps({"status": "registered"}), 200, {"Content-Type": "application/json"}
# ---------------------------------------------------------------------------
# Client provisioning API — used by Keenetic Unified (KU) to auto-create a
# unique Phobos client per router (name "ku-") and return its
# install command. Each call guarantees a unique WG keypair + tunnel IP +
# peer on ALL servers, so per-router configs never collide.
# ---------------------------------------------------------------------------
CLIENT_SCRIPT = f"{PHOBOS_DIR}/repo/server/scripts/phobos-client.sh"
def _normalize_client_id(name):
return name.strip().lower().replace(" ", "-")
def _latest_token_for_client(cid):
try:
with open(TOKENS_FILE) as f:
toks = json.load(f)
except Exception:
return ""
for t in reversed(toks):
if t.get("client") == cid:
return t.get("token", "")
return ""
@app.route("/api/client/ensure", methods=["POST"])
def api_client_ensure():
s = load_settings()
if request.headers.get("X-API-Key", "") != s.get("server_api_key", ""):
return json.dumps({"error": "unauthorized"}), 401, {"Content-Type": "application/json"}
data = request.json or {}
name = (data.get("name") or "").strip()
if not name or not re.match(r"^[a-zA-Z0-9_-]+$", name):
return json.dumps({"error": "bad_name"}), 400, {"Content-Type": "application/json"}
cid = _normalize_client_id(name)
env = {**os.environ, **_load_server_env()}
exists = Path(f"{CLIENTS_DIR}/{cid}").is_dir()
sub = "link" if exists else "add"
try:
out = subprocess.check_output([CLIENT_SCRIPT, sub, name], text=True,
timeout=60, stderr=subprocess.STDOUT, env=env)
except subprocess.CalledProcessError as e:
return json.dumps({"error": "script_failed", "output": (e.output or "")[-1500:]}), 500, {"Content-Type": "application/json"}
meta = next((c for c in get_clients() if c.get("client_id") == cid), None)
if not meta:
return json.dumps({"error": "no_meta", "output": out[-800:]}), 500, {"Content-Type": "application/json"}
pub = meta.get("public_key", "")
tip = (meta.get("tunnel_ip_v4") or "").split("/")[0]
# Peer on ALL servers so failover works on every server, not just primary.
if pub and tip:
try:
sync_peer_to_all_servers(pub, f"{tip}/32", "add")
except Exception:
pass
assigned = auto_assign_server(cid)
# Per-client pull token for the NAT-friendly config pull channel.
s = load_settings()
ra = s.get("router_access", {})
acc = ra.get(cid) or {}
if not acc.get("pull_token"):
acc["pull_token"] = secrets.token_hex(16)
ra[cid] = acc
s["router_access"] = ra
save_settings(s)
pull_token = ra[cid]["pull_token"]
token = _latest_token_for_client(cid)
# phobos-client.sh writes these 0600 (umask 077 in action_add) → nginx (www-data)
# can't read them → 403 on /init and /packages. Make them world-readable.
if token:
_www = f"{PHOBOS_DIR}/www"
for pth in (f"{_www}/init/{token}.sh",
f"{_www}/packages/{token}/phobos-{cid}.tar.gz"):
try:
os.chmod(pth, 0o644)
except Exception:
pass
try:
os.chmod(f"{_www}/packages/{token}", 0o755)
except Exception:
pass
install_url = f"http://{SERVER_IP}/init/{token}.sh" if token else ""
return json.dumps({
"ok": True, "client_id": cid, "install_url": install_url, "token": token,
"pull_token": pull_token, "tunnel_ip": tip, "assigned_server": assigned,
}), 200, {"Content-Type": "application/json"}
@app.route("/api/client/remove", methods=["POST"])
def api_client_remove():
s = load_settings()
if request.headers.get("X-API-Key", "") != s.get("server_api_key", ""):
return json.dumps({"error": "unauthorized"}), 401, {"Content-Type": "application/json"}
data = request.json or {}
cid = _normalize_client_id((data.get("name") or "").strip())
if not cid:
return json.dumps({"error": "bad_name"}), 400, {"Content-Type": "application/json"}
meta = next((c for c in get_clients() if c.get("client_id") == cid), None)
pub = (meta or {}).get("public_key", "")
tip = ((meta or {}).get("tunnel_ip_v4") or "").split("/")[0]
env = {**os.environ, **_load_server_env()}
try:
subprocess.check_output([CLIENT_SCRIPT, "remove", cid], text=True,
timeout=30, stderr=subprocess.STDOUT, env=env)
except subprocess.CalledProcessError as e:
return json.dumps({"error": "script_failed", "output": (e.output or "")[-1500:]}), 500, {"Content-Type": "application/json"}
if pub and tip:
try:
sync_peer_to_all_servers(pub, f"{tip}/32", "remove")
except Exception:
pass
s = load_settings()
for key in ("router_access", "client_assignments"):
d = s.get(key, {})
if cid in d:
del d[cid]
s[key] = d
save_settings(s)
return json.dumps({"ok": True, "client_id": cid}), 200, {"Content-Type": "application/json"}
@app.route("/api/obf-health")
def api_obf_health():
"""Liveness of the obfuscator path on THIS (primary) server.
Returns 200 only if all wg-obfuscator-* services are active, else 503.
Router check_primary uses this as a valid switchback signal (the panel
HTTP port stays up even when the tunnel path is dead, so it cannot be
used directly)."""
try:
listing = subprocess.check_output(
["systemctl", "list-units", "--type=service", "--no-legend",
"wg-obfuscator-*"], text=True, timeout=5)
names = [ln.split()[0] for ln in listing.splitlines() if ln.strip()]
except Exception:
names = []
if not names:
names = ["wg-obfuscator-2083.service",
"wg-obfuscator-5443.service",
"wg-obfuscator-993.service"]
bad = []
for n in names:
st = subprocess.getoutput("systemctl is-active " + n).strip()
if st != "active":
bad.append(n + "=" + st)
if bad:
return "obf-down: " + ",".join(bad), 503, {"Content-Type": "text/plain"}
return "obf-ok " + str(len(names)), 200, {"Content-Type": "text/plain"}
@app.route("/servers", methods=["GET", "POST"])
@auth_required
def servers_page():
s = load_settings()
servers = load_servers()
msg = ""
if request.method == "POST":
action = request.form.get("action")
if action == "add":
ip = request.form.get("ip", "").strip()
api_key = request.form.get("api_key", "").strip()
ssh_user = request.form.get("ssh_user", "root").strip() or "root"
ssh_pass = request.form.get("ssh_pass", "").strip()
if ip:
servers.append({
"ip": ip, "api_key": api_key, "enabled": True,
"last_seen": "", "wg_public_key": "", "obfuscator_key": "", "ports": "",
"ssh_user": ssh_user, "ssh_pass": ssh_pass
})
save_servers(servers)
order = s.get("server_order", [])
if ip not in order:
order.append(ip)
s["server_order"] = order
save_settings(s)
msg = f'Сервер {ip} добавлен'
elif action == "remove":
ip = request.form.get("ip", "").strip()
servers = [sv for sv in servers if sv.get("ip") != ip]
save_servers(servers)
order = s.get("server_order", [])
if ip in order:
order.remove(ip)
s["server_order"] = order
save_settings(s)
msg = f'Сервер удалён'
elif action == "sync":
ip = request.form.get("ip", "").strip()
srv = next((sv for sv in servers if sv["ip"] == ip), None)
if srv:
clients = get_clients()
synced = 0
for c in clients:
pub = c.get("public_key", "")
tip = c.get("tunnel_ip_v4", "")
if pub and tip:
res = sync_peer_to_server(srv, pub, f"{tip}/32", "add")
if res.get("status") == "ok":
synced += 1
msg = f'Синхронизировано {synced} клиентов на {ip}'
elif action == "fetch_info":
ip = request.form.get("ip", "").strip()
for srv in servers:
if srv["ip"] == ip:
try:
import urllib.request as ul
url = f"http://{ip}:8444/api/info"
rq = ul.Request(url, headers={"X-API-Key": srv.get("api_key", "")})
resp = ul.urlopen(rq, timeout=5)
info = json.loads(resp.read())
srv["wg_public_key"] = info.get("wg_public_key", "")
srv["obfuscator_key"] = info.get("obfuscator_key", "")
srv["ports"] = ",".join(info.get("ports", []))
save_servers(servers)
msg = f'Инфо получено от {ip}'
except Exception as e:
msg = f'Ошибка: {e}'
elif action == "set_api_key":
new_key = request.form.get("server_api_key", "").strip()
if new_key:
s["server_api_key"] = new_key
save_settings(s)
msg = 'API ключ обновлён'
elif action in ("move_up", "move_down"):
ip = request.form.get("ip", "").strip()
order = s.get("server_order", [])
if not order or SERVER_IP not in order:
order = [SERVER_IP] + [sv["ip"] for sv in servers]
if ip in order:
idx = order.index(ip)
if action == "move_up" and idx > 0:
order[idx], order[idx-1] = order[idx-1], order[idx]
elif action == "move_down" and idx < len(order) - 1:
order[idx], order[idx+1] = order[idx+1], order[idx]
s["server_order"] = order
save_settings(s)
elif action == "save_router_access":
clients = get_clients()
ra = s.get("router_access", {})
for c in clients:
cid = c.get("client_id", "")
ssh_ip = request.form.get(f"ssh_ip_{cid}", "").strip()
ssh_static = request.form.get(f"ssh_static_{cid}", "").strip()
ssh_user = request.form.get(f"ssh_user_{cid}", "root").strip() or "root"
ssh_pass = request.form.get(f"ssh_pass_{cid}", "").strip()
if ssh_ip or ssh_pass:
ra[cid] = {"ssh_ip": ssh_ip, "ssh_static": ssh_static, "ssh_user": ssh_user, "ssh_pass": ssh_pass, "ssh_ok": ra.get(cid, {}).get("ssh_ok", False)}
elif cid in ra and not ssh_ip and not ssh_pass:
del ra[cid]
s["router_access"] = ra
save_settings(s)
msg = 'SSH доступ сохранён'
elif action == "test_ssh":
cid = request.form.get("client_id", "").strip()
ra = s.get("router_access", {})
acc = ra.get(cid, {})
results = []
ok_ip = ""
# Try tunnel IP first, then static
for label, ip in [("tunnel", acc.get("ssh_ip", "")), ("static", acc.get("ssh_static", ""))]:
if not ip:
continue
try:
r = subprocess.run(
["sshpass", "-p", acc.get("ssh_pass", ""), "ssh",
"-o", "StrictHostKeyChecking=no", "-o", "ConnectTimeout=5",
f"{acc.get('ssh_user','root')}@{ip}", "echo OK"],
capture_output=True, text=True, timeout=10
)
if "OK" in r.stdout:
results.append(f"{label} ({ip}): ✓")
if not ok_ip:
ok_ip = ip
else:
results.append(f"{label} ({ip}): ✗ {r.stderr[:80]}")
except Exception as e:
results.append(f"{label} ({ip}): ✗ {str(e)[:80]}")
if cid in ra:
ra[cid]["ssh_ok"] = bool(ok_ip)
ra[cid]["ssh_tested"] = ok_ip or ""
s["router_access"] = ra
save_settings(s)
# SSH is optional under the pull model. Frame the result around whether
# the router is actually reachable/managed, not raw SSH success.
client_pub = ""
for _c in get_clients():
if _c.get("client_id") == cid:
client_pub = _c.get("public_key", "")
break
online = False
for _p, _i in get_wg_peers().items():
if _p == client_pub and is_peer_online(_i.get("latest handshake", "")):
online = True
for _ip, _hs in server_handshakes_cache.items():
_ts = _hs.get(client_pub, 0)
if _ts and (time.time() - _ts) < 180:
online = True
ssh_line = " | ".join(results) if results else tr("SSH не настроен", "SSH not configured")
if ok_ip:
msg = f'{cid}: SSH ✓ {ok_ip} — {ssh_line}'
elif online:
msg = f'{cid}: {tr("роутер ОНЛАЙН, управляется через pull. SSH по туннелю недоступен — это нормально для роутера за NAT (входящий SSH закрыт); управление SSH не требует.", "router is ONLINE, managed via pull. Tunnel SSH unreachable — normal for a NAT router (inbound SSH closed); management does not need SSH.")}'
else:
msg = f'{cid}: {tr("роутер НЕ виден — нет свежего handshake ни на одном сервере, и SSH недоступен.", "router NOT visible — no fresh handshake on any server, and SSH unreachable.")} {ssh_line}'
elif action == "push_all":
clients = get_clients()
results = []
for c in clients:
cid = c.get("client_id", "")
res = push_config_to_router(cid)
results.append(f"{cid}: {res['msg']}")
msg = '' + '
'.join(results) + ''
elif action == "save_server_access":
for srv in servers:
ip = srv["ip"]
su = request.form.get(f"srv_ssh_user_{ip}", "").strip()
sp = request.form.get(f"srv_ssh_pass_{ip}", "").strip()
if su:
srv["ssh_user"] = su
if sp:
srv["ssh_pass"] = sp
save_servers(servers)
msg = 'SSH доступ к серверам сохранён'
elif action == "rebalance":
clients = get_clients()
all_srv = get_all_servers_ordered()
server_ips = [sv["ip"] for sv in all_srv]
# Load-aware greedy: base load from live CPU+RAM (cached by monitor),
# then each client goes to the least-loaded server, with a per-client
# penalty so load spreads evenly. Down servers excluded (unless all down).
base_load = {ip: server_load(server_stats_cache.get(ip, {})) for ip in server_ips}
usable = [ip for ip in server_ips if base_load[ip] < 0.99] or server_ips
placed = {ip: 0 for ip in server_ips}
PEN = 0.08 # extra load each assigned client adds to a server's score
assignments = {}
rb_err = []
for c in clients:
cid = c["client_id"]
target = min(usable, key=lambda ip: base_load[ip] + placed[ip] * PEN)
placed[target] += 1
assignments[cid] = target
# Provision the WG peer on the target server. Round-robin assignment
# alone is not enough — without the peer the handshake fails there and
# the router's health monitor just fails back.
pub = c.get("public_key", "")
tip = c.get("tunnel_ip_v4", "").split("/")[0]
if pub and tip:
r = add_peer_to_server(target, pub, f"{tip}/32")
if r.get("status") != "ok":
rb_err.append(f"{cid}->{target}: {r.get('msg', 'peer fail')}")
s["client_assignments"] = assignments
save_settings(s)
for _cid in assignments:
fanout_router_config(_cid) # push conf to all servers (tunnel pull)
if rb_err:
msg = f'Rebalance: {"; ".join(rb_err)}'
else:
msg = f'{tr("Перераспределено", "Rebalanced")}: {len(clients)} -> {len(server_ips)} {tr("серв.; пиры добавлены, роутеры применят за ~15с", "servers; peers added, routers apply within ~15s")}'
# Build ordered server list
all_ordered = get_all_servers_ordered()
assignments = s.get("client_assignments", {})
clients = get_clients()
# Count assigned clients per server
assign_counts = {}
for cid, sip in assignments.items():
assign_counts[sip] = assign_counts.get(sip, 0) + 1
# Stats dashboard cards
stats_cards = ""
for srv in all_ordered:
ip = srv.get("ip", "")
is_primary = srv.get("is_primary", False)
cached = server_stats_cache.get(ip, {})
ports = srv.get("ports", "?")
assigned = assign_counts.get(ip, 0)
if is_primary:
local = get_local_stats()
cpu_val = local.get("cpu", "?")
mem_val = local.get("mem", "?")
peers_val = local.get("peers", 0)
st_class = "badge-on"
st_text = "Primary"
else:
cpu_val = cached.get("cpu", "—")
mem_val = cached.get("mem", "—")
peers_val = cached.get("peers", "?")
is_up = cached.get("status", "down") == "ok"
st_class = "badge-on" if is_up else "badge-off"
st_text = "Online" if is_up else "Offline"
stats_cards += f"""
{ip}
{st_text}
CPU: {cpu_val} RAM: {mem_val}
Peers: {peers_val} Assigned: {assigned} Ports: {ports}
"""
# Server priority rows
rows = ""
for idx, srv in enumerate(all_ordered):
ip = srv.get("ip", "")
is_primary = srv.get("is_primary", False)
ports = srv.get("ports", "?")
cached = server_stats_cache.get(ip, {})
peers_val = cached.get("peers", get_local_stats().get("peers", "?") if is_primary else "?")
assigned = assign_counts.get(ip, 0)
if is_primary:
st_class, st_text = "badge-on", "Primary"
else:
is_up = cached.get("status", "down") == "ok"
st_class = "badge-on" if is_up else "badge-off"
st_text = "Online" if is_up else "Offline"
prio = f'#{idx+1}'
ptag = ' Primary' if is_primary else ""
mv = f"""
"""
acts = ""
if not is_primary:
acts = f"""
"""
rows += f"""
{prio}{ip}{ptag} {ports} {peers_val} {assigned}
{st_text}
"""
api_key = s.get("server_api_key", "")
# Router access rows — tunnel IP + optional static IP
router_access = s.get("router_access", {})
access_rows = ""
for c in clients:
cid = c.get("client_id", "")
tunnel_ip = c.get("tunnel_ip_v4", "")
acc = router_access.get(cid, {})
ssh_ip = acc.get("ssh_ip", tunnel_ip)
ssh_static = acc.get("ssh_static", "")
ssh_user = acc.get("ssh_user", "root")
ssh_pass = acc.get("ssh_pass", "")
ssh_ok = acc.get("ssh_ok", False)
ssh_tested = acc.get("ssh_tested", "")
if ssh_ok:
badge_cls = "badge-on"
badge_txt = f"✓ {ssh_tested}" if ssh_tested else "✓"
elif ssh_ip and ssh_pass:
badge_cls = "badge-warn"
badge_txt = "не проверен"
else:
badge_cls = "badge-off"
badge_txt = "—"
access_rows += f"""
{cid}
{badge_txt}
"""
html = f"""
{nav('servers')}
{msg}
{tr("Панель серверов", "Server Dashboard")} {hlp("Живая статистика всех VPN-серверов. Peers = активные WG-подключения. Assigned = клиенты, у которых этот сервер основной. Фоновый монитор опрашивает каждые " + str(s.get('monitor_interval',30)) + " сек и шлёт оповещения в Telegram при падении/восстановлении сервера.", "Live stats for all VPN servers. Peers = active WG connections. Assigned = clients whose primary is this server. Background monitor polls every " + str(s.get('monitor_interval',30)) + "s and sends Telegram alerts on server up/down.")}
{stats_cards}
{tr("Приоритет серверов и балансировка", "Server Priority & Load Balancing")} {hlp("У каждого клиента есть основной сервер. При сбое роутер сам перебирает серверы сверху вниз. Стрелки меняют общий порядок резервирования. Роутеры за NAT подтягивают конфиг сами — ручных действий не нужно.", "Each client has a primary server. On failure the router tries servers top-to-bottom by itself. Arrows change the global fallback order. NAT routers pull config themselves — no manual action needed.")}
{tr("Сервер","Server")} {tr("Порты","Ports")} Peers {tr("Назначено","Assigned")} {tr("Статус","Status")} {tr("Действия","Actions")}
{rows}
{hlp("Равномерно распределить всех клиентов по серверам (round-robin). Меняет основной сервер у клиентов — они переключатся автоматически.", "Evenly redistribute all clients across servers (round-robin). Changes clients primary server — they switch automatically.")}
{hlp("Принудительно отправить конфиг на роутеры с прямым SSH. Роутеры за NAT это игнорируют и сами подтягивают конфиг за ~15 сек — обычно эта кнопка не нужна.", "Force-push config to routers reachable by SSH. NAT routers ignore it and auto-pull within ~15s — you usually do not need this.")}
{tr("SSH-доступ к роутерам (опционально)", "Router SSH Access (optional)")} {hlp("НЕ обязательно. Роутеры сами тянут конфиг с панели по HTTPS каждые ~15 сек (работает за любым NAT, входящий SSH не нужен). Эта секция — только для роутеров с прямым IP, если хочешь мгновенный push. У роутеров за NAT/KeenDNS тут будут таймауты SSH — это нормально и безвредно, pull-канал их обновляет.", "NOT required. Routers pull config from the panel over HTTPS every ~15s (works behind any NAT, no inbound SSH). This section is only for direct-IP routers if you want an instant push. NAT/KeenDNS routers will show SSH timeouts here — that is normal and harmless; the pull channel keeps them updated.")}
{"".join(f'' for c in clients)}
{tr("Команды для роутеров","Router commands")}
Установка — выполнить по SSH на роутере (Keenetic/Netcraze с Entware):
{_render_install_commands()}
Удаление Phobos — выполнить по SSH на роутере:
/opt/etc/Phobos/phobos-uninstall.sh
Скрипт остановит obfuscator, удалит cron, конфиги и бинарник. WireGuard интерфейс на роутере нужно удалить вручную через веб-панель.
Фикс SSH через туннель — если SSH по 10.25.0.x не работает (security-level):
wget -O - http://{SERVER_IP}/init/fix-security.sh | sh
Автоматически найдёт WG интерфейс Phobos и установит security-level private для входящих подключений (SSH). Выполнять на роутере по SSH через LAN (192.168.1.1).
Server SSH Access
SSH credentials for secondary servers. Used by panel to add/remove WG peers when switching client servers. Main server uses local commands.
API Key
Deploy Secondary Server
One command deploys WG + obfuscator + mini-API on a new VPS. Auto-registers in panel. Run via SSH on new VPS:
MAIN_SERVER={SERVER_IP} MAIN_API_KEY={api_key} bash <(curl -fsSL https://raw.githubusercontent.com/andrey271192/PCA_Phobos/main/server/secondary-setup.sh)
"""
return render(html)
def _load_server_env():
"""Load server.env as dict for subprocess env."""
env = {}
if os.path.exists(SERVER_ENV):
with open(SERVER_ENV) as f:
for line in f:
line = line.strip()
if "=" in line and not line.startswith("#"):
k, v = line.split("=", 1)
env[k] = v
return env
if __name__ == "__main__":
app.run(host="0.0.0.0", port=8443, debug=False)