Compare commits
11 commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6372d649da | |||
| 4791c5ca2b | |||
| a701d55e3b | |||
| e3bd279407 | |||
| 264991593a | |||
| 2c5ec8b06c | |||
| da6200ae4a | |||
| 6546d5bed1 | |||
| 695771c79a | |||
| 2b34b11391 | |||
| 02ff43095c |
18 changed files with 3462 additions and 310 deletions
221
.github/workflows/ci.yml
vendored
221
.github/workflows/ci.yml
vendored
|
|
@ -24,6 +24,10 @@ jobs:
|
|||
run: |
|
||||
bash -n install.sh
|
||||
bash -n mbs
|
||||
bash -n tests/test_mbs_update.sh
|
||||
|
||||
- name: Smoke test mbs update and mirror (manual edits on the server, url mirror, unsafe urls, rollback)
|
||||
run: bash tests/test_mbs_update.sh
|
||||
|
||||
- name: Smoke test install-script rendering
|
||||
env:
|
||||
|
|
@ -57,6 +61,9 @@ jobs:
|
|||
}
|
||||
script = nodeprov.render_install_script(node)
|
||||
assert "PRIVKEY" in script
|
||||
assert "PREFLIGHT_FAIL" in script and "443 2053 2087" in script, "node install must refuse a non-empty server before touching it"
|
||||
assert script.index("PREFLIGHT_FAIL") < script.index("authorized_keys"), "preflight must run before the management key is added"
|
||||
assert "OK=$((OK+1))" in script and "STATUS=failed" in script, "node must report active only after xray stays up"
|
||||
assert len(script) > 500
|
||||
print("node install script rendered OK,", len(script), "bytes")
|
||||
PYEOF
|
||||
|
|
@ -310,10 +317,10 @@ jobs:
|
|||
import settings
|
||||
|
||||
db.init_db()
|
||||
db.create_node("n1", "Node One", "managed", "1.1.1.1", 443, "pub1", "sid1", "sni1", "xtls-rprx-vision")
|
||||
db.create_node("bk1", "Node One", "managed", "1.1.1.1", 443, "pub1", "sid1", "sni1", "xtls-rprx-vision")
|
||||
|
||||
sub_a = db.create_subscription(111, "n1", 30, "1m", source="bot")
|
||||
sub_b = db.create_subscription(222, "n1", 30, "1m", source="bot")
|
||||
sub_a = db.create_subscription(111, "bk1", 30, "1m", source="bot")
|
||||
sub_b = db.create_subscription(222, "bk1", 30, "1m", source="bot")
|
||||
|
||||
legal.update_env_var("BRAND_NAME", "SnapshotBrand")
|
||||
settings.set_plan_prices({"1m": 555})
|
||||
|
|
@ -333,7 +340,7 @@ jobs:
|
|||
settings.set_plan_prices({"1m": 999})
|
||||
legal.update_env_var("HWID_LIMIT_ENABLED", "false")
|
||||
assert db.resume_subscription(sub_a["uuid"])["held_at"] is None
|
||||
sub_c = db.create_subscription(333, "n1", 30, "1m", source="bot")
|
||||
sub_c = db.create_subscription(333, "bk1", 30, "1m", source="bot")
|
||||
assert settings.get_brand_name() == "MutatedAfterBackup"
|
||||
assert len(db.list_active_subscriptions(tg_id=111)) == 1
|
||||
|
||||
|
|
@ -384,3 +391,209 @@ jobs:
|
|||
|
||||
print("custom ADMIN_PATH: old /admin route gone, new path registered, root() no longer leaks the panel OK")
|
||||
PYEOF
|
||||
|
||||
- name: Smoke test server chains (xray config generation, relay clients, subscription entries, audit log, old-db migration)
|
||||
env:
|
||||
BOT_TOKEN: "123456789:AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"
|
||||
BOT_USERNAME: "x"
|
||||
ADMIN_IDS: "1"
|
||||
ADMIN_PANEL_PASSWORD: "ci-test-password-not-real"
|
||||
PANEL_DOMAIN: "panel.test"
|
||||
SUB_DOMAIN: "sub.test"
|
||||
SITE_DOMAIN: "test"
|
||||
XRAY_PUBLIC_KEY: "x"
|
||||
XRAY_SHORT_ID_TCP: "x"
|
||||
XRAY_SHORT_ID_GRPC: "x"
|
||||
XRAY_SHORT_ID_XHTTP: "x"
|
||||
run: |
|
||||
python - << 'PYEOF'
|
||||
import base64
|
||||
import json
|
||||
import urllib.parse
|
||||
|
||||
import chains
|
||||
import db
|
||||
import links
|
||||
import nodeprov
|
||||
import xray_manager
|
||||
|
||||
transports = nodeprov.build_transports("a.example.com", 443, "www.microsoft.com", "PUBA", include_ws=False)
|
||||
entry_cfg = json.loads(nodeprov._build_config_json(transports, "PRIVA", "a.example.com"))
|
||||
exit_cfg = json.loads(nodeprov._build_config_json(
|
||||
nodeprov.build_transports("b.example.com", 443, "www.microsoft.com", "PUBB"), "PRIVB", "b.example.com"))
|
||||
|
||||
wanted = {"uuid-1": "uuid-1", "uuid-2": "uuid-2"}
|
||||
chain = {"code": "cabc12", "port": 10443, "short_id": "1234567890abcdef", "exit_node": "chb", "relay_uuid": "relay-uuid-1"}
|
||||
exit_nodes = {"chb": {
|
||||
"address": "b.example.com", "port": 443, "sni": "www.microsoft.com",
|
||||
"public_key": "PUBB", "short_id": "ffff", "kind": "managed", "shared_uuid": None,
|
||||
}}
|
||||
|
||||
base_before = json.dumps([ib for ib in entry_cfg["inbounds"] if ib["tag"] in chains.BASE_TAGS], sort_keys=True)
|
||||
|
||||
changed, problems = chains.sync_config(entry_cfg, wanted, {}, [chain], exit_nodes)
|
||||
assert changed and not problems, problems
|
||||
tags = [ib["tag"] for ib in entry_cfg["inbounds"]]
|
||||
assert "chain-cabc12" in tags, tags
|
||||
ci = chains.find_inbound(entry_cfg, "chain-cabc12")
|
||||
assert ci["port"] == 10443
|
||||
assert ci["streamSettings"]["realitySettings"]["shortIds"] == ["1234567890abcdef"]
|
||||
assert ci["streamSettings"]["realitySettings"]["privateKey"] == "PRIVA"
|
||||
assert [c["id"] for c in ci["settings"]["clients"]] == ["uuid-1", "uuid-2"]
|
||||
assert all(c["flow"] == "xtls-rprx-vision" for c in ci["settings"]["clients"])
|
||||
out = [o for o in entry_cfg["outbounds"] if o["tag"] == "chain-cabc12-out"]
|
||||
assert len(out) == 1
|
||||
vn = out[0]["settings"]["vnext"][0]
|
||||
assert vn["address"] == "b.example.com" and vn["port"] == 443 and vn["users"][0]["id"] == "relay-uuid-1"
|
||||
assert out[0]["streamSettings"]["realitySettings"]["publicKey"] == "PUBB"
|
||||
rules = [r for r in entry_cfg["routing"]["rules"] if r.get("outboundTag") == "chain-cabc12-out"]
|
||||
assert len(rules) == 1 and rules[0]["inboundTag"] == ["chain-cabc12"]
|
||||
assert entry_cfg["routing"]["rules"][0]["outboundTag"] == "api"
|
||||
assert entry_cfg["outbounds"][0]["tag"] == "direct", "default outbound must stay first"
|
||||
print("entry config: chain inbound/outbound/rule built OK")
|
||||
|
||||
changed2, problems2 = chains.sync_config(entry_cfg, wanted, {}, [chain], exit_nodes)
|
||||
assert not changed2 and not problems2, "second pass must be a no-op"
|
||||
print("idempotent OK")
|
||||
|
||||
wanted3 = {"uuid-2": "uuid-2", "uuid-3": "uuid-3"}
|
||||
changed3, _ = chains.sync_config(entry_cfg, wanted3, {}, [chain], exit_nodes)
|
||||
assert changed3
|
||||
ci = chains.find_inbound(entry_cfg, "chain-cabc12")
|
||||
assert [c["id"] for c in ci["settings"]["clients"]] == ["uuid-2", "uuid-3"]
|
||||
for tag in ("vless-tcp-reality", "vless-grpc-reality", "vless-xhttp-reality"):
|
||||
assert [c["id"] for c in chains.find_inbound(entry_cfg, tag)["settings"]["clients"]] == ["uuid-2", "uuid-3"]
|
||||
print("clients follow the active set on every inbound incl. chain OK")
|
||||
|
||||
changed4, _ = chains.sync_config(entry_cfg, {"uuid-9": "uuid-9"}, {}, [], exit_nodes, apply_chains=False)
|
||||
assert changed4
|
||||
assert chains.find_inbound(entry_cfg, "chain-cabc12") is not None, "apply_chains=False must not drop chains"
|
||||
assert [c["id"] for c in chains.find_inbound(entry_cfg, "chain-cabc12")["settings"]["clients"]] == ["uuid-9"]
|
||||
print("clients-only fallback keeps existing chains and still syncs their clients OK")
|
||||
|
||||
chains.sync_config(entry_cfg, wanted, {}, [], exit_nodes)
|
||||
assert chains.find_inbound(entry_cfg, "chain-cabc12") is None
|
||||
assert not [o for o in entry_cfg["outbounds"] if o["tag"].startswith("chain-")]
|
||||
assert not [r for r in entry_cfg["routing"]["rules"] if str(r.get("outboundTag", "")).startswith("chain-")]
|
||||
base_after = json.dumps([ib for ib in entry_cfg["inbounds"] if ib["tag"] in chains.BASE_TAGS], sort_keys=True)
|
||||
assert json.loads(base_after) != [] and len(json.loads(base_after)) == len(json.loads(base_before))
|
||||
print("removing the chain cleans inbound/outbound/rule OK")
|
||||
|
||||
changed5, p5 = chains.sync_config(exit_cfg, wanted, {"relay-uuid-1": "relay-cabc12"}, [], {})
|
||||
assert changed5 and not p5
|
||||
tcp_ids = [c["id"] for c in chains.find_inbound(exit_cfg, "vless-tcp-reality")["settings"]["clients"]]
|
||||
grpc_ids = [c["id"] for c in chains.find_inbound(exit_cfg, "vless-grpc-reality")["settings"]["clients"]]
|
||||
assert "relay-uuid-1" in tcp_ids and "relay-uuid-1" not in grpc_ids
|
||||
relay_entry = [c for c in chains.find_inbound(exit_cfg, "vless-tcp-reality")["settings"]["clients"] if c["id"] == "relay-uuid-1"][0]
|
||||
assert relay_entry["flow"] == "xtls-rprx-vision" and relay_entry["email"] == "relay-cabc12"
|
||||
changed6, _ = chains.sync_config(exit_cfg, wanted, {"relay-uuid-1": "relay-cabc12"}, [], {})
|
||||
assert not changed6
|
||||
print("exit node keeps the relay client only on the TCP inbound and survives sync OK")
|
||||
|
||||
busy_cfg = json.loads(nodeprov._build_config_json(transports, "PRIVA", "a.example.com"))
|
||||
usable, skipped = chains.split_busy_chains(busy_cfg, [chain], {10443})
|
||||
assert usable == [] and len(skipped) == 1
|
||||
usable2, skipped2 = chains.split_busy_chains(busy_cfg, [chain], set())
|
||||
assert usable2 == [chain] and skipped2 == []
|
||||
chains.sync_config(busy_cfg, wanted, {}, [chain], exit_nodes)
|
||||
usable3, skipped3 = chains.split_busy_chains(busy_cfg, [chain], {10443})
|
||||
assert usable3 == [chain], "an already-applied chain is not a new port, busy check must ignore it"
|
||||
print("busy port handling OK")
|
||||
|
||||
ext_nodes = {"chb": dict(exit_nodes["chb"], kind="external", shared_uuid="shared-1")}
|
||||
ext_chain = dict(chain, relay_uuid=None)
|
||||
cfg_e = json.loads(nodeprov._build_config_json(transports, "PRIVA", "a.example.com"))
|
||||
ch, pr = chains.sync_config(cfg_e, wanted, {}, [ext_chain], ext_nodes)
|
||||
assert ch and not pr
|
||||
assert [o for o in cfg_e["outbounds"] if o["tag"] == "chain-cabc12-out"][0]["settings"]["vnext"][0]["users"][0]["id"] == "shared-1"
|
||||
no_key = dict(ext_nodes["chb"], shared_uuid=None)
|
||||
cfg_f = json.loads(nodeprov._build_config_json(transports, "PRIVA", "a.example.com"))
|
||||
ch, pr = chains.sync_config(cfg_f, wanted, {}, [ext_chain], {"chb": no_key})
|
||||
assert pr and chains.find_inbound(cfg_f, "chain-cabc12") is None
|
||||
print("external exit uses shared uuid, missing key is reported OK")
|
||||
|
||||
assert chains.latency_level(10) == "low" and chains.latency_level(80) == "medium" and chains.latency_level(300) == "high"
|
||||
assert chains.latency_level(None) == "unknown"
|
||||
assert chains.median_ms([-1, -1]) is None and chains.median_ms([30, 10, -1]) == 30
|
||||
print("latency helpers OK")
|
||||
|
||||
db.init_db()
|
||||
db.create_node("cha", "🇫🇮 Финляндия", "managed", "fi.example.com", 443, "PUBFI", "sidfi", "www.microsoft.com", "xtls-rprx-vision")
|
||||
db.create_node("chb", "🇳🇱 Нидерланды", "managed", "nl.example.com", 443, "PUBNL", "sidnl", "www.microsoft.com", "xtls-rprx-vision")
|
||||
db.create_node("chx", "Внешняя", "external", "ex.example.com", 443, "PUBEX", "sidex", "www.microsoft.com", "xtls-rprx-vision", shared_uuid="shared-ex")
|
||||
|
||||
c1 = db.create_chain("Финка → Голландия", "cha", "chb", "relay-1")
|
||||
assert c1["port"] == 10443 and len(c1["short_id"]) == 16 and c1["code"].startswith("c")
|
||||
c2 = db.create_chain("Финка → Внешняя", "cha", "chx", None)
|
||||
assert c2["port"] == 10444
|
||||
try:
|
||||
db.create_chain("dup", "cha", "chb", "x")
|
||||
assert False
|
||||
except ValueError:
|
||||
pass
|
||||
try:
|
||||
db.delete_node("chb")
|
||||
assert False, "node used in chain must not be deletable"
|
||||
except ValueError as e:
|
||||
assert "chain" in str(e)
|
||||
assert [c["code"] for c in db.list_chains()] == [c1["code"], c2["code"]]
|
||||
assert len(db.list_chains(enabled_only=True)) == 2
|
||||
db.update_chain(c2["code"], enabled=0)
|
||||
assert len(db.list_chains(enabled_only=True)) == 1
|
||||
assert db.stats()["chains"] == 1
|
||||
print("db chains CRUD, port allocation, node-delete guard OK")
|
||||
|
||||
sub = db.create_subscription(500, "cha", 30, "1m", source="bot")
|
||||
text = base64.b64decode(links.build_subscription_text([sub])).decode()
|
||||
lines = text.split("\n")
|
||||
chain_lines = [l for l in lines if ":10443?" in l]
|
||||
assert len(chain_lines) == 1, lines
|
||||
assert "10444" not in text, "disabled chain must not leak into the subscription"
|
||||
parsed = urllib.parse.urlparse(chain_lines[0])
|
||||
assert parsed.hostname == "fi.example.com" and parsed.port == 10443
|
||||
qs = urllib.parse.parse_qs(parsed.query)
|
||||
assert qs["sid"] == [c1["short_id"]] and qs["pbk"] == ["PUBFI"] and qs["flow"] == ["xtls-rprx-vision"]
|
||||
assert urllib.parse.unquote(parsed.fragment) == "🇫🇮 Финляндия → 🇳🇱 Нидерланды"
|
||||
assert parsed.username == sub["uuid"]
|
||||
print("subscription text carries the chain entry for the entry node's subscribers OK")
|
||||
|
||||
other = db.create_subscription(501, "chb", 30, "1m", source="bot")
|
||||
text2 = base64.b64decode(links.build_subscription_text([other])).decode()
|
||||
assert ":10443?" not in text2, "subscribers of the exit node must not get the entry node's chain"
|
||||
print("chain is only offered to entry-node subscribers OK")
|
||||
|
||||
db.update_node("cha", enabled=0)
|
||||
text3 = base64.b64decode(links.build_subscription_text([sub])).decode()
|
||||
assert text3.strip() == ""
|
||||
db.update_node("cha", enabled=1)
|
||||
|
||||
db.update_chain(c2["code"], enabled=1)
|
||||
node_n1 = db.get_node("cha")
|
||||
w, relay, entry_chains, exit_n = xray_manager.desired_state(node_n1)
|
||||
assert sub["uuid"] in w and [c["code"] for c in entry_chains] == [c1["code"], c2["code"]] and relay == {}
|
||||
node_n2 = db.get_node("chb")
|
||||
w2, relay2, entry2, exit2 = xray_manager.desired_state(node_n2)
|
||||
assert relay2 == {"relay-1": chains.relay_email(c1["code"])} and entry2 == []
|
||||
node_ex = db.get_node("chx")
|
||||
w3, relay3, entry3, exit3 = xray_manager.desired_state(node_ex)
|
||||
assert relay3 == {}
|
||||
db.update_node("chb", enabled=0)
|
||||
w4, relay4, entry4, exit4 = xray_manager.desired_state(node_n1)
|
||||
assert [c["code"] for c in entry4] == [c2["code"]], "chain whose exit is disabled must drop out"
|
||||
print("desired_state: entry/relay/disabled-node logic OK")
|
||||
|
||||
db.add_audit("admin", "node.add", "/admin/api/nodes", "1.2.3.4")
|
||||
db.add_audit(None, "login.failed", "", "5.6.7.8")
|
||||
rows = db.list_audit(10)
|
||||
assert rows[0]["action"] == "login.failed" and rows[1]["admin"] == "admin"
|
||||
print("audit log OK")
|
||||
|
||||
with db.get_conn() as conn:
|
||||
conn.execute("DROP TABLE chains")
|
||||
conn.execute("DROP TABLE audit_log")
|
||||
db.init_db()
|
||||
assert db.list_chains() == [] and db.list_audit() == []
|
||||
print("init_db recreates chain/audit tables on an old database OK")
|
||||
|
||||
print("chains: all smoke tests passed")
|
||||
PYEOF
|
||||
|
|
|
|||
2
.gitignore
vendored
2
.gitignore
vendored
|
|
@ -7,3 +7,5 @@ __pycache__/
|
|||
*.pyc
|
||||
venv/
|
||||
.claude/
|
||||
.update_mirror
|
||||
local-changes/
|
||||
|
|
|
|||
44
README.md
44
README.md
|
|
@ -1,8 +1,8 @@
|
|||
# MBS Panel
|
||||
|
||||
[](https://github.com/devsavsis/mbs-panel/actions/workflows/ci.yml)
|
||||
[](https://github.com/savsisbtw/mbs-panel/actions/workflows/ci.yml)
|
||||
[](LICENSE)
|
||||
[](https://github.com/devsavsis/mbs-panel/releases)
|
||||
[](https://github.com/savsisbtw/mbs-panel/releases)
|
||||
[](https://www.python.org/)
|
||||
[](https://github.com/XTLS/Xray-core)
|
||||
|
||||
|
|
@ -30,8 +30,12 @@
|
|||
- **Rate-limit на вход** — по IP, отдельно на пароль и на 2FA-код.
|
||||
- **Свой путь входа** — страницу логина можно увести с дефолтного `/admin` на любой другой (`ADMIN_PATH` в `.env`), доп. слой поверх rate-limit и 2FA — у Remnawave это в списке заявленных мер безопасности, у Marzban нет вообще.
|
||||
- **Пауза подписки** — временно отключить доступ без потери оплаченных дней (Marzban это умеет, Remnawave — нет): при возобновлении срок сдвигается ровно на длительность паузы.
|
||||
- **Исходящие вебхуки** — на оплату, выдачу/отзыв/паузу/возобновление подписки и на добавление/удаление/вкл-выкл ноды, с HMAC-подписью тела. У Remnawave это события по юзерам и нодам, у Marzban — только по юзерам; мы покрываем оба класса.
|
||||
- **Поиск и фильтр по подпискам** — по юзернейму/tg id/ноде/тарифу и по статусу, прямо в таблице.
|
||||
- **Исходящие вебхуки** — на оплату, выдачу/отзыв/паузу/возобновление подписки и на добавление/удаление/вкл-выкл ноды и создание/удаление/вкл-выкл цепочки, с HMAC-подписью тела. У Remnawave это события по юзерам и нодам, у Marzban — только по юзерам; мы покрываем оба класса.
|
||||
- **Поиск и фильтр по подпискам** — по юзернейму/tg id/ноде/тарифу и по статусу, прямо в таблице. Плюс экспорт всех подписок в CSV одной кнопкой.
|
||||
- **Цепочки серверов (v1.6)** — `клиент → нода A → нода B → интернет`: собираются в админке мышкой, зажал ЛКМ на клиенте и протянул провод через серверы к интернету. Больше двух серверов нельзя специально, на двух панель предупреждает про задержку и сама меряет RTT между нодами. Подробности ниже в разделе «Цепочки серверов».
|
||||
- **Юзеры и выдача подписок (v1.6)** — отдельная страница со всеми, кто уже есть в системе, даже если подписок у них ни разу не было (в «Подписках» таких не видно): поиск по юзернейму и Telegram ID, счётчики активных подписок и устройств, кнопка «Выдать подписку». Можно выдать и по Telegram ID тому, кого в базе ещё нет, подписка дождётся, пока он зайдёт в бота.
|
||||
- **Журнал действий (v1.6)** — кто из админов и когда менял ноды, цепочки, подписки, настройки, качал бэкап и входил (включая неудачные входы с IP). Пароли и ключи в журнал не попадают, только факт действия.
|
||||
- **Пинг нод и палитра команд (v1.6)** — живая задержка от панели до каждой ноды на дашборде и в списке нод; `Ctrl K` открывает поиск по страницам, нодам, цепочкам и действиям. Акцентный цвет панели меняется кружками сверху.
|
||||
|
||||
## Архитектура
|
||||
|
||||
|
|
@ -111,7 +115,7 @@ sequenceDiagram
|
|||
bash <(curl -Ls https://mbs.savsis.xyz/install.sh)
|
||||
```
|
||||
|
||||
(или напрямую с GitHub, если так удобнее: `git clone https://github.com/devsavsis/mbs-panel.git && cd mbs-panel && sudo bash install.sh` — скрипт один и тот же, `mbs.savsis.xyz` просто зеркало с автосинком)
|
||||
(или напрямую с GitHub, если так удобнее: `git clone https://github.com/savsisbtw/mbs-panel.git && cd mbs-panel && sudo bash install.sh` — скрипт один и тот же, `mbs.savsis.xyz` просто зеркало с автосинком)
|
||||
|
||||
Скрипт спросит домен панели, домен подписки, токен бота от [@BotFather](https://t.me/BotFather) и список Telegram ID админов — и дальше всё сам: ставит зависимости, Xray, nginx, выпускает сертификаты Let's Encrypt, генерирует Reality-ключи, поднимает systemd-сервисы, настраивает firewall (ufw) и fail2ban. В конце покажет пароль от админки и ссылку на панель.
|
||||
|
||||
|
|
@ -148,10 +152,18 @@ mbs status статус bot / api / xray / nginx
|
|||
mbs restart перезапустить bot + api
|
||||
mbs logs [bot|api|xray] последние строки лога (по умолчанию api)
|
||||
mbs domain текущий домен панели
|
||||
mbs update обновить код с GitHub и перезапустить
|
||||
mbs backup полная копия панели в /root/mbs-backups (база, .env, твои правки, конфиг Xray)
|
||||
mbs update [ссылка] обновить код и перезапустить; со ссылкой на git-зеркало берёт обновление оттуда
|
||||
mbs mirror [ссылка|off] показать / запомнить / убрать своё зеркало, его mbs update проверяет первым
|
||||
```
|
||||
|
||||
`mbs update` тянет `git pull` (только fast-forward — если на сервере что-то правили руками, честно откажется и не полезет мержить), ставит зависимости, **проверяет, что новый код вообще компилируется**, и только потом перезапускает. Если после рестарта `mbs-bot`/`mbs-api` не поднялись — сам откатывает на предыдущий коммит и поднимает его. `.env` и база (`mbs.db`) не в гите — их не тронет ни при каком раскладе.
|
||||
`mbs update` тянет обновление (только fast-forward, чужую историю на сервере не мержит), ставит зависимости, **проверяет, что новый код вообще компилируется**, и только потом перезапускает. Если после рестарта `mbs-bot`/`mbs-api` не поднялись — сам откатывает на предыдущий коммит и поднимает его. `.env` и база (`mbs.db`) не в гите — их не тронет ни при каком раскладе.
|
||||
|
||||
Перед каждым обновлением `mbs update` сам снимает полную копию в `/root/mbs-backups/` (консистентный снапшот базы через backup API SQLite, `.env`, код вместе с твоими ручными правками, конфиг Xray, без `venv`), хранит 5 последних, права `600`. Не получилось сделать копию (нет места на диске), обновление даже не начнётся. То же вручную: `mbs backup`. Вернуть всё как было: `tar xzf /root/mbs-backups/mbs-before-update-<время>.tar.gz -C /` и `mbs restart`.
|
||||
|
||||
Если на сервере правили файлы руками (бывает, `bot.py`/`config.py`/`db.py` под себя), обновление больше на этом не падает: правки откладываются в `git stash` и сохраняются патчем в `local-changes/local-changes-<время>.patch`, потом подтягивается новая версия. Вернуть своё поверх новой: `git stash pop` (может быть конфликт, если новая версия правила те же строки, тогда смотри патч). Если новый код не прошёл проверку или сервисы не поднялись, откат на старый коммит возвращает и твои правки.
|
||||
|
||||
Источники по порядку: своё зеркало (если задано через `mbs mirror`), потом `api.savsis.xyz`, потом GitHub. Появилось новое зеркало или GitHub недоступен, а ссылка на репо есть: `mbs update https://example.com/путь/mbs-panel.git` возьмёт обновление именно оттуда, один раз. Чтобы всегда обновляться с него: `mbs mirror https://example.com/путь/mbs-panel.git` (убрать: `mbs mirror off`). Принимаются только `https://`, `http://`, `ssh://` и `git@хост:путь`, всё остальное (в том числе `file://` и хитрые транспорты типа `ext::`) отбрасывается, ветка берётся `main`.
|
||||
|
||||
## Добавление ноды
|
||||
|
||||
|
|
@ -178,6 +190,22 @@ sequenceDiagram
|
|||
Panel->>Panel: нода активна, доступна в боте
|
||||
```
|
||||
|
||||
## Цепочки серверов
|
||||
|
||||
Обычное подключение это `клиент → нода → интернет`. Цепочка добавляет второй прыжок: `клиент → нода A → нода B → интернет`. Сайты видят IP ноды B, а клиент коннектится к A. Пригождается, когда вход хочется держать в одном регионе (ближе, не режут), а выход нужен в другой стране. Больше двух серверов не даёт специально: каждый лишний прыжок это задержка, а скорость упирается в самое слабое звено.
|
||||
|
||||
Собирается в админке: Цепочки → зажимаешь ЛКМ на «Клиенте», тянешь провод через серверы из пула и отпускаешь на «Интернете». Пока ведёшь, сервер под курсором цепляется после короткой задержки (чтоб не хватать всё подряд по пути). Можно и без перетаскивания, просто кликами по серверам и по «Интернету», `Esc` сбрасывает. Один сервер это обычное подключение, оно и так есть у каждой ноды, а вот два уже цепочка: панель сразу показывает предупреждение про высокую задержку и замеряет реальный RTT между нодами (TCP-коннект с входной ноды до выходной).
|
||||
|
||||
Что реально происходит под капотом:
|
||||
|
||||
- на входной ноде появляется отдельный Xray-inbound `chain-<код>` (тот же Reality-ключ что у ноды, но свой порт из 10443–10999 и свой shortId), outbound `chain-<код>-out` до выходной ноды и routing-правило «всё из этого inbound уходит в этот outbound»;
|
||||
- на выходной ноде заводится служебный клиент `relay-<код>`, под ним входная нода и ходит на выход (только на TCP+Reality inbound, на обоих прыжках `xtls-rprx-vision`);
|
||||
- порт открывается в ufw сам, конфиг прогоняется через `xray run -test`, после рестарта проверяется что Xray реально поднялся, если нет, конфиг откатывается;
|
||||
- подписчикам входной ноды в подписку добавляется ещё одна ссылка «A → B», подписчикам выходной цепочку не выдаём;
|
||||
- любая проблема с цепочкой не блокирует обычную синхронизацию клиентов, они применятся в любом случае.
|
||||
|
||||
Ограничения: входом может быть только локальная или управляемая нода (панель правит её конфиг), выходом ещё и внешняя нода с общим UUID. Ноду, которая сидит в цепочке, удалить нельзя, сначала удали цепочку. И важный момент про Reality: SNI-маскировка (`dest`) не должна быть сайтом с пост-квантовым обменом ключами (например `www.microsoft.com`), на таком Reality не заводится вообще, ни в цепочке, ни без неё, проверено руками. Дефолтный `www.wildberries.ru` подходит.
|
||||
|
||||
## Приём оплаты
|
||||
|
||||
По умолчанию бот выдаёт подписки бесплатно по кнопке — платежи выключены (`PAYMENTS_ENABLED=false`). Чтобы продавать доступ:
|
||||
|
|
@ -215,7 +243,7 @@ sequenceDiagram
|
|||
PR и issues welcome. CI на каждый пуш гоняет compile-check по питону, синтаксис-проверку шелл-скриптов и smoke-тест генерации install-скрипта ноды.
|
||||
|
||||
## Авторы:
|
||||
github.com/devsavsis
|
||||
github.com/savsisbtw
|
||||
github.com/welfizx
|
||||
|
||||
## Лицензия
|
||||
|
|
|
|||
1935
admin.html
1935
admin.html
File diff suppressed because it is too large
Load diff
291
api.py
291
api.py
|
|
@ -1,17 +1,22 @@
|
|||
import asyncio
|
||||
import csv
|
||||
import datetime
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import secrets
|
||||
import subprocess
|
||||
import urllib.request
|
||||
import uuid as uuidlib
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from fastapi import FastAPI, HTTPException, Request, Response, File, UploadFile
|
||||
from fastapi.middleware.cors import CORSMiddleware
|
||||
from fastapi import Body
|
||||
from fastapi.responses import HTMLResponse, PlainTextResponse, FileResponse
|
||||
|
||||
import backup
|
||||
import chains
|
||||
import db
|
||||
import legal
|
||||
import links
|
||||
|
|
@ -52,6 +57,59 @@ def require_admin(request: Request):
|
|||
raise HTTPException(401, "unauthorized")
|
||||
|
||||
|
||||
AUDIT_RULES = [
|
||||
("POST", r"^/admin/api/nodes$", "node.add"),
|
||||
("POST", r"^/admin/api/nodes/provision-guide$", "node.provision"),
|
||||
("POST", r"^/admin/api/nodes/reorder$", "node.reorder"),
|
||||
("PATCH", r"^/admin/api/nodes/[^/]+$", "node.edit"),
|
||||
("DELETE", r"^/admin/api/nodes/[^/]+$", "node.delete"),
|
||||
("POST", r"^/admin/api/chains$", "chain.create"),
|
||||
("PATCH", r"^/admin/api/chains/[^/]+$", "chain.edit"),
|
||||
("DELETE", r"^/admin/api/chains/[^/]+$", "chain.delete"),
|
||||
("POST", r"^/admin/api/subscriptions/[^/]+/revoke$", "sub.revoke"),
|
||||
("POST", r"^/admin/api/subscriptions/[^/]+/hold$", "sub.hold"),
|
||||
("POST", r"^/admin/api/subscriptions/[^/]+/resume$", "sub.resume"),
|
||||
("POST", r"^/admin/api/subscriptions/[^/]+/reset-traffic$", "sub.reset_traffic"),
|
||||
("POST", r"^/admin/api/users/-?\d+/grant$", "user.grant"),
|
||||
("POST", r"^/admin/api/users/-?\d+/hwid-limit$", "user.hwid_limit"),
|
||||
("DELETE", r"^/admin/api/users/-?\d+/devices/\d+$", "user.device_delete"),
|
||||
("POST", r"^/admin/api/gift-codes$", "gift.create"),
|
||||
("POST", r"^/admin/api/admins$", "admin.add"),
|
||||
("DELETE", r"^/admin/api/admins/\d+$", "admin.delete"),
|
||||
("POST", r"^/admin/api/2fa/enable$", "2fa.enable"),
|
||||
("POST", r"^/admin/api/2fa/disable$", "2fa.disable"),
|
||||
("GET", r"^/admin/api/backup$", "backup.download"),
|
||||
("POST", r"^/admin/api/backup/restore$", "backup.restore"),
|
||||
("POST", r"^/admin/api/settings/bot$", "settings.bot"),
|
||||
("POST", r"^/admin/api/branding$", "settings.brand"),
|
||||
("POST", r"^/admin/api/webhook-settings$", "settings.webhook"),
|
||||
("POST", r"^/admin/api/hwid-settings$", "settings.hwid"),
|
||||
("POST", r"^/admin/api/payments/(yookassa-settings|platega-settings|plan-settings|legal-settings)$", "settings.payments"),
|
||||
]
|
||||
AUDIT_COMPILED = [(method, re.compile(pattern), action) for method, pattern, action in AUDIT_RULES]
|
||||
|
||||
|
||||
def _audit_action(method: str, path: str):
|
||||
for rule_method, pattern, action in AUDIT_COMPILED:
|
||||
if rule_method == method and pattern.match(path):
|
||||
return action
|
||||
return None
|
||||
|
||||
|
||||
@app.middleware("http")
|
||||
async def audit_middleware(request: Request, call_next):
|
||||
path = request.url.path
|
||||
action = _audit_action(request.method, path) if path.startswith("/admin/api/") else None
|
||||
admin_name = None
|
||||
if action:
|
||||
admin = await asyncio.to_thread(db.get_session_admin, request.cookies.get(ADMIN_COOKIE))
|
||||
admin_name = admin["username"] if admin else None
|
||||
response = await call_next(request)
|
||||
if action and response.status_code < 400:
|
||||
await asyncio.to_thread(db.add_audit, admin_name, action, path, _client_ip(request))
|
||||
return response
|
||||
|
||||
|
||||
def _days_left(sub: dict) -> int:
|
||||
exp = datetime.datetime.fromisoformat(sub["expires_at"])
|
||||
reference = datetime.datetime.fromisoformat(sub["held_at"]) if sub.get("held_at") else datetime.datetime.utcnow()
|
||||
|
|
@ -614,11 +672,13 @@ def admin_login(request: Request, response: Response, body: dict = Body(...)):
|
|||
admin = db.verify_admin_login(username, password)
|
||||
if not admin:
|
||||
db.record_login_attempt(ip, "password")
|
||||
db.add_audit(username[:64] or None, "login.failed", "", ip)
|
||||
raise HTTPException(401, "wrong username or password")
|
||||
db.clear_login_attempts(ip, "password")
|
||||
if admin.get("totp_secret"):
|
||||
pending_token = db.create_pending_totp(admin["id"])
|
||||
return {"ok": True, "needs_totp": True, "pending_token": pending_token}
|
||||
db.add_audit(admin["username"], "login.ok", "", ip)
|
||||
token = db.create_admin_session(admin["id"])
|
||||
response.set_cookie(ADMIN_COOKIE, token, httponly=True, secure=True, samesite="strict", max_age=7 * 24 * 3600)
|
||||
return {"ok": True}
|
||||
|
|
@ -637,9 +697,11 @@ def admin_login_totp(request: Request, response: Response, body: dict = Body(...
|
|||
admin = db.get_admin_by_id(pending["admin_id"])
|
||||
if not admin or not admin.get("totp_secret") or not totp.verify(admin["totp_secret"], code):
|
||||
db.record_login_attempt(ip, "totp")
|
||||
db.add_audit(admin["username"] if admin else None, "login.totp_failed", "", ip)
|
||||
raise HTTPException(401, "wrong code")
|
||||
db.clear_login_attempts(ip, "totp")
|
||||
db.delete_pending_totp(pending_token)
|
||||
db.add_audit(admin["username"], "login.ok", "2fa", ip)
|
||||
token = db.create_admin_session(admin["id"])
|
||||
response.set_cookie(ADMIN_COOKIE, token, httponly=True, secure=True, samesite="strict", max_age=7 * 24 * 3600)
|
||||
return {"ok": True}
|
||||
|
|
@ -942,12 +1004,22 @@ def admin_reset_traffic(uuid: str, request: Request):
|
|||
return {"ok": ok}
|
||||
|
||||
|
||||
@app.get("/admin/api/users")
|
||||
def admin_list_users(request: Request, q: str = "", limit: int = 200):
|
||||
require_admin(request)
|
||||
return db.list_users(q=q, limit=limit)
|
||||
|
||||
|
||||
@app.get("/admin/api/users/{tg_id}")
|
||||
def admin_user_card(tg_id: int, request: Request):
|
||||
require_admin(request)
|
||||
user = db.get_user(tg_id)
|
||||
if not user:
|
||||
raise HTTPException(404, "not found")
|
||||
return {
|
||||
"tg_id": tg_id, "username": None, "created_at": None, "token": None, "exists": False,
|
||||
"subscriptions": [], "devices": [], "hwid_limit": None,
|
||||
"hwid_fallback_limit": settings.get_hwid_settings()["fallback_limit"],
|
||||
}
|
||||
subs = db.list_subscriptions_for_user(tg_id)
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
plans_by_code = settings.get_plans_by_code()
|
||||
|
|
@ -966,6 +1038,7 @@ def admin_user_card(tg_id: int, request: Request):
|
|||
"username": user["username"],
|
||||
"created_at": user["created_at"],
|
||||
"token": user["token"],
|
||||
"exists": True,
|
||||
"subscriptions": out_subs,
|
||||
"devices": db.list_devices(tg_id),
|
||||
"hwid_limit": user.get("hwid_limit"),
|
||||
|
|
@ -1231,6 +1304,222 @@ def admin_node_status(code: str, request: Request):
|
|||
return {"status": node["status"], "enabled": bool(node["enabled"])}
|
||||
|
||||
|
||||
def _panel_latency(node: dict) -> dict:
|
||||
samples = chains.tcp_connect_ms(node["address"], node["port"], samples=2, timeout=2.0)
|
||||
return {"code": node["code"], "ms": chains.median_ms(samples)}
|
||||
|
||||
|
||||
@app.get("/admin/api/nodes/latency")
|
||||
def admin_nodes_latency(request: Request):
|
||||
require_admin(request)
|
||||
targets = [n for n in db.list_nodes() if n["address"] and n["status"] == "active" and n["enabled"]]
|
||||
if not targets:
|
||||
return []
|
||||
with ThreadPoolExecutor(max_workers=min(8, len(targets))) as pool:
|
||||
return list(pool.map(_panel_latency, targets))
|
||||
|
||||
|
||||
def _chain_view(chain: dict, nodes_by_code: dict) -> dict:
|
||||
entry = nodes_by_code.get(chain["entry_node"])
|
||||
exit_node = nodes_by_code.get(chain["exit_node"])
|
||||
return {
|
||||
"code": chain["code"], "label": chain["label"],
|
||||
"entry_node": chain["entry_node"], "exit_node": chain["exit_node"],
|
||||
"entry_label": entry["label"] if entry else chain["entry_node"],
|
||||
"exit_label": exit_node["label"] if exit_node else chain["exit_node"],
|
||||
"entry_address": entry["address"] if entry else None,
|
||||
"port": chain["port"], "enabled": bool(chain["enabled"]),
|
||||
"created_at": chain["created_at"],
|
||||
}
|
||||
|
||||
|
||||
def _csv_safe(value) -> str:
|
||||
text = str(value)
|
||||
if text and text[0] in "=+-@\t\r":
|
||||
return "'" + text
|
||||
return text
|
||||
|
||||
|
||||
def _sync_chain_nodes(chain: dict) -> dict:
|
||||
results = {}
|
||||
for code in dict.fromkeys([chain["entry_node"], chain["exit_node"]]):
|
||||
node = db.get_node(code)
|
||||
if not node:
|
||||
continue
|
||||
try:
|
||||
res = xray_manager.sync_node(node)
|
||||
results[code] = {"ok": not res["problems"], "problems": res["problems"]}
|
||||
except Exception as e:
|
||||
results[code] = {"ok": False, "problems": [str(e)]}
|
||||
return results
|
||||
|
||||
|
||||
def _chain_failures(results: dict) -> list:
|
||||
failures = []
|
||||
for code, res in results.items():
|
||||
for problem in res["problems"]:
|
||||
failures.append(f"{code}: {problem}")
|
||||
return failures
|
||||
|
||||
|
||||
@app.get("/admin/api/chains")
|
||||
def admin_list_chains(request: Request):
|
||||
require_admin(request)
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
return [_chain_view(c, nodes_by_code) for c in db.list_chains()]
|
||||
|
||||
|
||||
@app.post("/admin/api/chains")
|
||||
def admin_create_chain(request: Request, body: dict = Body(...)):
|
||||
require_admin(request)
|
||||
entry_code = str(body.get("entry", "")).strip()
|
||||
exit_code = str(body.get("exit", "")).strip()
|
||||
entry = db.get_node(entry_code)
|
||||
exit_node = db.get_node(exit_code)
|
||||
if not entry or not exit_node:
|
||||
raise HTTPException(400, "выбери входной и выходной серверы из списка нод")
|
||||
if entry_code == exit_code:
|
||||
raise HTTPException(400, "вход и выход цепочки должны быть разными серверами")
|
||||
if entry["kind"] not in chains.CHAIN_KINDS_ENTRY:
|
||||
raise HTTPException(400, "входной сервер должен быть под управлением панели (локальный или управляемый)")
|
||||
if exit_node["kind"] not in chains.CHAIN_KINDS_EXIT:
|
||||
raise HTTPException(400, "этот тип ноды нельзя использовать как выход цепочки")
|
||||
for node in (entry, exit_node):
|
||||
if not node["enabled"] or node["status"] != "active":
|
||||
raise HTTPException(400, f"нода «{node['label']}» выключена или ещё не установлена")
|
||||
relay_uuid = None
|
||||
if exit_node["kind"] == "external":
|
||||
if not exit_node.get("shared_uuid"):
|
||||
raise HTTPException(400, "у внешней ноды не задан shared UUID — через неё цепочку не построить")
|
||||
else:
|
||||
relay_uuid = str(uuidlib.uuid4())
|
||||
label = str(body.get("label", "")).strip()[:80] or f"{entry['label']} → {exit_node['label']}"
|
||||
try:
|
||||
chain = db.create_chain(label, entry_code, exit_code, relay_uuid)
|
||||
except ValueError as e:
|
||||
raise HTTPException(400, str(e))
|
||||
results = _sync_chain_nodes(chain)
|
||||
failures = _chain_failures(results)
|
||||
if failures:
|
||||
db.delete_chain(chain["code"])
|
||||
_sync_chain_nodes(chain)
|
||||
raise HTTPException(502, "цепочка не применилась, всё откатили назад: " + "; ".join(failures))
|
||||
webhooks.send("chain.created", {
|
||||
"code": chain["code"], "label": chain["label"], "entry": entry_code, "exit": exit_code,
|
||||
})
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
return _chain_view(chain, nodes_by_code)
|
||||
|
||||
|
||||
@app.patch("/admin/api/chains/{code}")
|
||||
def admin_update_chain(code: str, request: Request, body: dict = Body(...)):
|
||||
require_admin(request)
|
||||
chain = db.get_chain(code)
|
||||
if not chain:
|
||||
raise HTTPException(404, "not found")
|
||||
fields = {}
|
||||
if "label" in body:
|
||||
label = str(body["label"]).strip()[:80]
|
||||
if not label:
|
||||
raise HTTPException(400, "название не может быть пустым")
|
||||
fields["label"] = label
|
||||
toggled = "enabled" in body and bool(body["enabled"]) != bool(chain["enabled"])
|
||||
if "enabled" in body:
|
||||
fields["enabled"] = 1 if body["enabled"] else 0
|
||||
updated = db.update_chain(code, **fields)
|
||||
if toggled:
|
||||
failures = _chain_failures(_sync_chain_nodes(updated))
|
||||
if failures:
|
||||
db.update_chain(code, enabled=chain["enabled"])
|
||||
_sync_chain_nodes(chain)
|
||||
raise HTTPException(502, "не получилось применить: " + "; ".join(failures))
|
||||
webhooks.send("chain.enabled" if updated["enabled"] else "chain.disabled", {
|
||||
"code": code, "label": updated["label"],
|
||||
})
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
return _chain_view(db.get_chain(code), nodes_by_code)
|
||||
|
||||
|
||||
@app.delete("/admin/api/chains/{code}")
|
||||
def admin_delete_chain(code: str, request: Request):
|
||||
require_admin(request)
|
||||
chain = db.get_chain(code)
|
||||
if not chain:
|
||||
raise HTTPException(404, "not found")
|
||||
db.delete_chain(code)
|
||||
results = _sync_chain_nodes(chain)
|
||||
webhooks.send("chain.deleted", {"code": code, "label": chain["label"]})
|
||||
return {"ok": True, "warnings": _chain_failures(results)}
|
||||
|
||||
|
||||
def _hop_probe(entry: dict, exit_node: dict) -> dict:
|
||||
try:
|
||||
samples = xray_manager.probe_from_node(entry, exit_node["address"], exit_node["port"])
|
||||
except Exception as e:
|
||||
return {"rtt_ms": None, "samples": [], "level": "unknown", "error": str(e)}
|
||||
rtt = chains.median_ms(samples)
|
||||
return {"rtt_ms": rtt, "samples": samples, "level": chains.latency_level(rtt)}
|
||||
|
||||
|
||||
@app.get("/admin/api/chains/probe")
|
||||
def admin_probe_chain(entry: str, exit: str, request: Request):
|
||||
require_admin(request)
|
||||
entry_node = db.get_node(entry)
|
||||
exit_node = db.get_node(exit)
|
||||
if not entry_node or not exit_node or entry == exit:
|
||||
raise HTTPException(400, "нужны две разные ноды")
|
||||
return _hop_probe(entry_node, exit_node)
|
||||
|
||||
|
||||
@app.post("/admin/api/chains/{code}/check")
|
||||
def admin_check_chain(code: str, request: Request):
|
||||
require_admin(request)
|
||||
chain = db.get_chain(code)
|
||||
if not chain:
|
||||
raise HTTPException(404, "not found")
|
||||
entry = db.get_node(chain["entry_node"])
|
||||
exit_node = db.get_node(chain["exit_node"])
|
||||
if not entry or not exit_node:
|
||||
raise HTTPException(400, "одна из нод цепочки удалена")
|
||||
entry_alive = nodeprov.check_node_alive(entry["address"], chain["port"])
|
||||
hop = _hop_probe(entry, exit_node)
|
||||
return {"entry_alive": entry_alive, **hop}
|
||||
|
||||
|
||||
@app.get("/admin/api/audit")
|
||||
def admin_audit(request: Request, limit: int = 100):
|
||||
require_admin(request)
|
||||
return db.list_audit(limit=max(1, min(limit, 500)))
|
||||
|
||||
|
||||
@app.get("/admin/api/subscriptions/export.csv")
|
||||
def admin_export_subscriptions(request: Request):
|
||||
require_admin(request)
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
plans_by_code = settings.get_plans_by_code()
|
||||
buf = io.StringIO()
|
||||
writer = csv.writer(buf)
|
||||
writer.writerow(["tg_id", "username", "node", "plan", "created_at", "expires_at", "days_left", "status", "uuid"])
|
||||
for s in db.list_all_subscriptions(limit=100000):
|
||||
node = nodes_by_code.get(s["node"])
|
||||
plan = plans_by_code.get(s["plan"])
|
||||
if s.get("held_at"):
|
||||
status = "held"
|
||||
elif not s["active"] or s["expires_at"] <= db.now_iso():
|
||||
status = "expired"
|
||||
else:
|
||||
status = "active"
|
||||
writer.writerow([_csv_safe(v) for v in [
|
||||
s["tg_id"], s.get("username") or "", node["label"] if node else s["node"],
|
||||
plan["label"] if plan else s["plan"], s["created_at"], s["expires_at"],
|
||||
_days_left(s), status, s["uuid"],
|
||||
]])
|
||||
return Response(
|
||||
content="" + buf.getvalue(), media_type="text/csv; charset=utf-8",
|
||||
headers={"Content-Disposition": 'attachment; filename="subscriptions.csv"'},
|
||||
)
|
||||
|
||||
|
||||
|
||||
ADMIN_HTML_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "admin.html")
|
||||
|
||||
|
|
|
|||
|
|
@ -123,5 +123,8 @@ def restore_backup(data: bytes) -> dict:
|
|||
os.chmod(tmp_db_path, 0o600)
|
||||
os.replace(tmp_db_path, DB_PATH)
|
||||
|
||||
import db
|
||||
db.init_db()
|
||||
|
||||
_prune_old_safety_copies()
|
||||
return {"restored_env": restored_env, "safety_copy": safety_copy}
|
||||
|
|
|
|||
41
bot.py
41
bot.py
|
|
@ -33,6 +33,7 @@ def main_menu_kb(tg_id: int) -> InlineKeyboardMarkup:
|
|||
rows = [
|
||||
[InlineKeyboardButton(text="Получить VPN", callback_data="menu:get")],
|
||||
[InlineKeyboardButton(text="Моя подписка", callback_data="menu:mysub")],
|
||||
[InlineKeyboardButton(text="Пригласить друга", callback_data="menu:referral")],
|
||||
[InlineKeyboardButton(text="О сервисе", callback_data="menu:about")],
|
||||
]
|
||||
if is_admin(tg_id):
|
||||
|
|
@ -40,6 +41,14 @@ def main_menu_kb(tg_id: int) -> InlineKeyboardMarkup:
|
|||
return InlineKeyboardMarkup(inline_keyboard=rows)
|
||||
|
||||
|
||||
async def get_bot_username() -> str:
|
||||
global _bot_username
|
||||
if _bot_username is None:
|
||||
me = await bot.get_me()
|
||||
_bot_username = me.username
|
||||
return _bot_username
|
||||
|
||||
|
||||
def nodes_kb(prefix: str) -> InlineKeyboardMarkup:
|
||||
rows = []
|
||||
for n in db.list_nodes(enabled_only=True):
|
||||
|
|
@ -90,6 +99,12 @@ async def send_main_menu(message: Message):
|
|||
async def start_deeplink(message: Message, command: CommandObject):
|
||||
user = db.get_or_create_user(message.from_user.id, message.from_user.username)
|
||||
payload = command.args or ""
|
||||
if payload.startswith("ref_") or payload.startswith("ref-"):
|
||||
ref_code = payload[4:]
|
||||
referrer = db.get_user_by_ref_code(ref_code)
|
||||
if referrer and settings.get_referral_settings()["enabled"]:
|
||||
db.set_referred_by(message.from_user.id, referrer["tg_id"])
|
||||
return await send_main_menu(message)
|
||||
if payload.startswith("gift_") or payload.startswith("gift-"):
|
||||
code = payload[5:]
|
||||
gift, err = db.redeem_gift_code(code, message.from_user.id)
|
||||
|
|
@ -140,6 +155,32 @@ async def cb_about(cb: CallbackQuery):
|
|||
await cb.answer()
|
||||
|
||||
|
||||
@dp.callback_query(F.data == "menu:referral")
|
||||
async def cb_referral(cb: CallbackQuery):
|
||||
kb = InlineKeyboardMarkup(inline_keyboard=[[InlineKeyboardButton(text="Назад", callback_data="menu:main")]])
|
||||
ref_settings = settings.get_referral_settings()
|
||||
if not ref_settings["enabled"]:
|
||||
await cb.message.edit_text("Реферальная программа сейчас отключена.", reply_markup=kb)
|
||||
return await cb.answer()
|
||||
user = db.get_or_create_user(cb.from_user.id, cb.from_user.username)
|
||||
stats = db.referral_stats(cb.from_user.id)
|
||||
username = await get_bot_username()
|
||||
link = f"https://t.me/{username}?start=ref_{user['ref_code']}"
|
||||
days = ref_settings["bonus_days"]
|
||||
text = (
|
||||
f"<b>Пригласи друга</b>\n\n"
|
||||
f"За каждого друга, который активирует подписку по твоей ссылке, "
|
||||
f"вы <b>оба</b> получаете +{days} дн. к подписке.\n\n"
|
||||
f"{DIVIDER}\n"
|
||||
f"Твоя ссылка:\n<code>{link}</code>\n\n"
|
||||
f"Приглашено: {stats['referred_count']}\n"
|
||||
)
|
||||
if stats["bonus_days_pending"]:
|
||||
text += f"Накоплено бонусных дней (зачислятся при следующей подписке): {stats['bonus_days_pending']}\n"
|
||||
await cb.message.edit_text(text, reply_markup=kb)
|
||||
await cb.answer()
|
||||
|
||||
|
||||
@dp.callback_query(F.data == "menu:get")
|
||||
async def cb_get(cb: CallbackQuery):
|
||||
await cb.message.edit_text("Выбери сервер:", reply_markup=nodes_kb("node"))
|
||||
|
|
|
|||
233
chains.py
Normal file
233
chains.py
Normal file
|
|
@ -0,0 +1,233 @@
|
|||
import copy
|
||||
import json
|
||||
import re
|
||||
import socket
|
||||
import time
|
||||
|
||||
BASE_TAGS = ("vless-tcp-reality", "vless-grpc-reality", "vless-xhttp-reality", "vless-ws-tls")
|
||||
TCP_TAG = "vless-tcp-reality"
|
||||
VISION = "xtls-rprx-vision"
|
||||
CHAIN_PREFIX = "chain-"
|
||||
RELAY_EMAIL_PREFIX = "relay-"
|
||||
MAX_SERVERS = 2
|
||||
PORT_MIN = 10443
|
||||
PORT_MAX = 10999
|
||||
CODE_RE = re.compile(r"^[a-z0-9]{1,16}$")
|
||||
HOST_RE = re.compile(r"^[A-Za-z0-9.-]{1,253}$")
|
||||
CHAIN_KINDS_ENTRY = ("local", "managed")
|
||||
CHAIN_KINDS_EXIT = ("local", "managed", "external")
|
||||
|
||||
|
||||
class ChainConfigError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
def inbound_tag(code):
|
||||
return CHAIN_PREFIX + code
|
||||
|
||||
|
||||
def outbound_tag(code):
|
||||
return CHAIN_PREFIX + code + "-out"
|
||||
|
||||
|
||||
def relay_email(code):
|
||||
return RELAY_EMAIL_PREFIX + code
|
||||
|
||||
|
||||
def is_chain_inbound_tag(tag):
|
||||
return bool(tag) and tag.startswith(CHAIN_PREFIX) and not tag.endswith("-out")
|
||||
|
||||
|
||||
def is_user_tag(tag):
|
||||
return tag in BASE_TAGS or is_chain_inbound_tag(tag)
|
||||
|
||||
|
||||
def flow_for_tag(tag):
|
||||
if tag == TCP_TAG or is_chain_inbound_tag(tag):
|
||||
return VISION
|
||||
return None
|
||||
|
||||
|
||||
def sync_clients(clients, wanted, flow):
|
||||
kept = []
|
||||
seen = set()
|
||||
for c in clients:
|
||||
cid = c.get("id")
|
||||
if cid in wanted and cid not in seen:
|
||||
kept.append(c)
|
||||
seen.add(cid)
|
||||
for cid in wanted:
|
||||
if cid in seen:
|
||||
continue
|
||||
entry = {"id": cid, "email": wanted[cid]}
|
||||
if flow:
|
||||
entry["flow"] = flow
|
||||
kept.append(entry)
|
||||
return kept
|
||||
|
||||
|
||||
def find_inbound(cfg, tag):
|
||||
for ib in cfg.get("inbounds", []):
|
||||
if ib.get("tag") == tag:
|
||||
return ib
|
||||
return None
|
||||
|
||||
|
||||
def build_chain_inbound(template, chain, wanted, old_clients):
|
||||
reality = (template.get("streamSettings") or {}).get("realitySettings")
|
||||
if not reality:
|
||||
raise ChainConfigError("у входной ноды нет TCP+Reality inbound — цепочку строить не из чего")
|
||||
ib = copy.deepcopy(template)
|
||||
ib["tag"] = inbound_tag(chain["code"])
|
||||
ib["port"] = chain["port"]
|
||||
ib["streamSettings"]["realitySettings"]["shortIds"] = [chain["short_id"]]
|
||||
ib["settings"]["clients"] = sync_clients(old_clients, wanted, VISION)
|
||||
return ib
|
||||
|
||||
|
||||
def build_chain_outbound(chain, exit_node, relay_uuid):
|
||||
return {
|
||||
"tag": outbound_tag(chain["code"]),
|
||||
"protocol": "vless",
|
||||
"settings": {
|
||||
"vnext": [{
|
||||
"address": exit_node["address"],
|
||||
"port": int(exit_node["port"]),
|
||||
"users": [{"id": relay_uuid, "encryption": "none", "flow": VISION}],
|
||||
}],
|
||||
},
|
||||
"streamSettings": {
|
||||
"network": "tcp",
|
||||
"security": "reality",
|
||||
"realitySettings": {
|
||||
"serverName": exit_node["sni"],
|
||||
"fingerprint": "chrome",
|
||||
"publicKey": exit_node["public_key"],
|
||||
"shortId": exit_node["short_id"],
|
||||
"spiderX": "",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def build_chain_rule(chain):
|
||||
return {
|
||||
"type": "field",
|
||||
"inboundTag": [inbound_tag(chain["code"])],
|
||||
"outboundTag": outbound_tag(chain["code"]),
|
||||
}
|
||||
|
||||
|
||||
def relay_for_chain(chain, exit_node):
|
||||
if exit_node["kind"] == "external":
|
||||
return exit_node.get("shared_uuid")
|
||||
return chain.get("relay_uuid")
|
||||
|
||||
|
||||
def split_busy_chains(cfg, entry_chains, busy_ports):
|
||||
new_ports = set(new_ports_needed(cfg, entry_chains))
|
||||
usable = []
|
||||
problems = []
|
||||
for chain in entry_chains:
|
||||
if chain["port"] in new_ports and chain["port"] in busy_ports:
|
||||
problems.append(f"{chain['code']}: порт {chain['port']} уже занят другим процессом, цепочка не применена")
|
||||
continue
|
||||
usable.append(chain)
|
||||
return usable, problems
|
||||
|
||||
|
||||
def sync_config(cfg, wanted, relay_wanted, entry_chains, exit_nodes, apply_chains=True):
|
||||
before = json.dumps(cfg, sort_keys=True)
|
||||
problems = []
|
||||
|
||||
template = find_inbound(cfg, TCP_TAG)
|
||||
|
||||
for ib in cfg["inbounds"]:
|
||||
tag = ib.get("tag")
|
||||
if not is_user_tag(tag):
|
||||
continue
|
||||
want = dict(wanted)
|
||||
if tag == TCP_TAG:
|
||||
want.update(relay_wanted)
|
||||
ib["settings"]["clients"] = sync_clients(ib["settings"]["clients"], want, flow_for_tag(tag))
|
||||
|
||||
if not apply_chains:
|
||||
changed = json.dumps(cfg, sort_keys=True) != before
|
||||
return changed, problems
|
||||
|
||||
old_chain_inbounds = {}
|
||||
for ib in cfg["inbounds"]:
|
||||
if is_chain_inbound_tag(ib.get("tag")):
|
||||
old_chain_inbounds[ib["tag"]] = ib
|
||||
|
||||
kept_inbounds = [ib for ib in cfg["inbounds"] if not is_chain_inbound_tag(ib.get("tag"))]
|
||||
kept_outbounds = [ob for ob in cfg.get("outbounds", []) if not (ob.get("tag") or "").startswith(CHAIN_PREFIX)]
|
||||
routing = cfg.setdefault("routing", {})
|
||||
kept_rules = [r for r in routing.get("rules", []) if not (r.get("outboundTag") or "").startswith(CHAIN_PREFIX)]
|
||||
|
||||
for chain in entry_chains:
|
||||
exit_node = exit_nodes.get(chain["exit_node"])
|
||||
if template is None:
|
||||
problems.append(f"{chain['code']}: нет TCP+Reality inbound на входной ноде")
|
||||
continue
|
||||
if not exit_node:
|
||||
problems.append(f"{chain['code']}: выходная нода не найдена")
|
||||
continue
|
||||
relay_uuid = relay_for_chain(chain, exit_node)
|
||||
if not relay_uuid:
|
||||
problems.append(f"{chain['code']}: у выходной ноды нет ключа для цепочки")
|
||||
continue
|
||||
old = old_chain_inbounds.get(inbound_tag(chain["code"]))
|
||||
old_clients = old["settings"]["clients"] if old else []
|
||||
kept_inbounds.append(build_chain_inbound(template, chain, wanted, old_clients))
|
||||
kept_outbounds.append(build_chain_outbound(chain, exit_node, relay_uuid))
|
||||
kept_rules.append(build_chain_rule(chain))
|
||||
|
||||
cfg["inbounds"] = kept_inbounds
|
||||
cfg["outbounds"] = kept_outbounds
|
||||
routing["rules"] = kept_rules
|
||||
|
||||
changed = json.dumps(cfg, sort_keys=True) != before
|
||||
return changed, problems
|
||||
|
||||
|
||||
def new_ports_needed(cfg, entry_chains):
|
||||
existing = set()
|
||||
for ib in cfg.get("inbounds", []):
|
||||
if is_chain_inbound_tag(ib.get("tag")):
|
||||
existing.add(ib["tag"])
|
||||
ports = []
|
||||
for chain in entry_chains:
|
||||
if inbound_tag(chain["code"]) not in existing:
|
||||
ports.append(chain["port"])
|
||||
return ports
|
||||
|
||||
|
||||
def tcp_connect_ms(host, port, samples=3, timeout=3.0):
|
||||
results = []
|
||||
for _ in range(samples):
|
||||
start = time.perf_counter()
|
||||
try:
|
||||
with socket.create_connection((host, int(port)), timeout=timeout):
|
||||
pass
|
||||
results.append(round((time.perf_counter() - start) * 1000))
|
||||
except OSError:
|
||||
results.append(-1)
|
||||
return results
|
||||
|
||||
|
||||
def median_ms(samples):
|
||||
good = sorted(s for s in samples if s >= 0)
|
||||
if not good:
|
||||
return None
|
||||
return good[len(good) // 2]
|
||||
|
||||
|
||||
def latency_level(rtt_ms):
|
||||
if rtt_ms is None:
|
||||
return "unknown"
|
||||
if rtt_ms < 40:
|
||||
return "low"
|
||||
if rtt_ms < 120:
|
||||
return "medium"
|
||||
return "high"
|
||||
|
|
@ -120,3 +120,6 @@ PLATEGA_SECRET = env("PLATEGA_SECRET", "")
|
|||
|
||||
HWID_LIMIT_ENABLED = env("HWID_LIMIT_ENABLED", "false").lower() == "true"
|
||||
HWID_FALLBACK_LIMIT = int(env("HWID_FALLBACK_LIMIT", "3"))
|
||||
|
||||
REFERRAL_ENABLED = env("REFERRAL_ENABLED", "true").lower() == "true"
|
||||
REFERRAL_BONUS_DAYS = int(env("REFERRAL_BONUS_DAYS", "3"))
|
||||
|
|
|
|||
235
db.py
235
db.py
|
|
@ -6,6 +6,7 @@ import secrets
|
|||
import datetime
|
||||
import contextlib
|
||||
|
||||
import chains as chainsmod
|
||||
from config import DB_PATH
|
||||
|
||||
SCHEMA = """
|
||||
|
|
@ -113,6 +114,31 @@ CREATE TABLE IF NOT EXISTS nodes (
|
|||
created_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS chains (
|
||||
code TEXT PRIMARY KEY,
|
||||
label TEXT NOT NULL,
|
||||
entry_node TEXT NOT NULL,
|
||||
exit_node TEXT NOT NULL,
|
||||
port INTEGER NOT NULL,
|
||||
short_id TEXT NOT NULL,
|
||||
relay_uuid TEXT,
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
sort_order INTEGER NOT NULL DEFAULT 0,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS audit_log (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
ts TEXT NOT NULL,
|
||||
admin TEXT,
|
||||
action TEXT NOT NULL,
|
||||
detail TEXT,
|
||||
ip TEXT
|
||||
);
|
||||
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_chains_pair ON chains (entry_node, exit_node);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_chains_entry_port ON chains (entry_node, port);
|
||||
CREATE INDEX IF NOT EXISTS idx_audit_ts ON audit_log (ts);
|
||||
CREATE INDEX IF NOT EXISTS idx_subs_active_expires ON subscriptions (active, expires_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_subs_tg_id ON subscriptions (tg_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_subs_node ON subscriptions (node);
|
||||
|
|
@ -133,6 +159,10 @@ _NEW_NODE_COLUMNS = {
|
|||
|
||||
_NEW_USER_COLUMNS = {
|
||||
"hwid_limit": "INTEGER",
|
||||
"ref_code": "TEXT",
|
||||
"referred_by": "INTEGER",
|
||||
"referral_rewarded": "INTEGER NOT NULL DEFAULT 0",
|
||||
"bonus_days_pending": "INTEGER NOT NULL DEFAULT 0",
|
||||
}
|
||||
|
||||
_NEW_ADMIN_SESSION_COLUMNS = {
|
||||
|
|
@ -159,6 +189,7 @@ def _migrate():
|
|||
for name, decl in _NEW_USER_COLUMNS.items():
|
||||
if name not in ucols:
|
||||
conn.execute(f"ALTER TABLE users ADD COLUMN {name} {decl}")
|
||||
conn.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_users_ref_code ON users(ref_code)")
|
||||
scols = {r["name"] for r in conn.execute("PRAGMA table_info(admin_sessions)").fetchall()}
|
||||
for name, decl in _NEW_ADMIN_SESSION_COLUMNS.items():
|
||||
if name not in scols:
|
||||
|
|
@ -350,25 +381,178 @@ def delete_node(code: str):
|
|||
).fetchone()["c"]
|
||||
if active:
|
||||
raise ValueError(f"node has {active} active subscriptions, revoke them first")
|
||||
in_chains = conn.execute(
|
||||
"SELECT COUNT(*) c FROM chains WHERE entry_node=? OR exit_node=?", (code, code)
|
||||
).fetchone()["c"]
|
||||
if in_chains:
|
||||
raise ValueError(f"node is used in {in_chains} chain(s), delete them first")
|
||||
conn.execute("DELETE FROM nodes WHERE code=?", (code,))
|
||||
|
||||
|
||||
def list_chains(enabled_only: bool = False):
|
||||
q = "SELECT * FROM chains"
|
||||
if enabled_only:
|
||||
q += " WHERE enabled=1"
|
||||
q += " ORDER BY sort_order ASC, created_at ASC"
|
||||
with get_conn() as conn:
|
||||
rows = conn.execute(q).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def get_chain(code: str):
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT * FROM chains WHERE code=?", (code,)).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
|
||||
def create_chain(label: str, entry_node: str, exit_node: str, relay_uuid: str | None):
|
||||
with get_conn() as conn:
|
||||
dup = conn.execute(
|
||||
"SELECT 1 FROM chains WHERE entry_node=? AND exit_node=?", (entry_node, exit_node)
|
||||
).fetchone()
|
||||
if dup:
|
||||
raise ValueError("такая цепочка уже есть")
|
||||
used = {r["port"] for r in conn.execute(
|
||||
"SELECT port FROM chains WHERE entry_node=?", (entry_node,)
|
||||
).fetchall()}
|
||||
port = None
|
||||
for candidate in range(chainsmod.PORT_MIN, chainsmod.PORT_MAX + 1):
|
||||
if candidate not in used:
|
||||
port = candidate
|
||||
break
|
||||
if port is None:
|
||||
raise ValueError("закончились свободные порты под цепочки на этой ноде")
|
||||
row = conn.execute("SELECT MAX(sort_order) m FROM chains").fetchone()
|
||||
next_order = (row["m"] or 0) + 1
|
||||
code = "c" + secrets.token_hex(3)
|
||||
conn.execute(
|
||||
"INSERT INTO chains (code, label, entry_node, exit_node, port, short_id, relay_uuid, enabled, sort_order, created_at) "
|
||||
"VALUES (?,?,?,?,?,?,?,1,?,?)",
|
||||
(code, label, entry_node, exit_node, port, secrets.token_hex(8), relay_uuid, next_order, now_iso()),
|
||||
)
|
||||
return get_chain(code)
|
||||
|
||||
|
||||
def update_chain(code: str, **fields):
|
||||
if not fields:
|
||||
return get_chain(code)
|
||||
cols = ", ".join(f"{k}=?" for k in fields)
|
||||
with get_conn() as conn:
|
||||
conn.execute(f"UPDATE chains SET {cols} WHERE code=?", (*fields.values(), code))
|
||||
return get_chain(code)
|
||||
|
||||
|
||||
def delete_chain(code: str):
|
||||
with get_conn() as conn:
|
||||
conn.execute("DELETE FROM chains WHERE code=?", (code,))
|
||||
|
||||
|
||||
def add_audit(admin: str | None, action: str, detail: str = "", ip: str | None = None):
|
||||
with get_conn() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO audit_log (ts, admin, action, detail, ip) VALUES (?,?,?,?,?)",
|
||||
(now_iso(), admin, action, detail[:500], ip),
|
||||
)
|
||||
conn.execute("DELETE FROM audit_log WHERE id <= (SELECT MAX(id) FROM audit_log) - 5000")
|
||||
|
||||
|
||||
def list_audit(limit: int = 100):
|
||||
with get_conn() as conn:
|
||||
rows = conn.execute("SELECT * FROM audit_log ORDER BY id DESC LIMIT ?", (limit,)).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def _generate_ref_code(conn) -> str:
|
||||
for _ in range(20):
|
||||
code = secrets.token_hex(4)
|
||||
if not conn.execute("SELECT 1 FROM users WHERE ref_code=?", (code,)).fetchone():
|
||||
return code
|
||||
raise RuntimeError("could not generate a unique ref_code")
|
||||
|
||||
|
||||
def get_or_create_user(tg_id: int, username: str | None):
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
if row:
|
||||
if username and row["username"] != username:
|
||||
conn.execute("UPDATE users SET username=? WHERE tg_id=?", (username, tg_id))
|
||||
if not row["ref_code"]:
|
||||
conn.execute(
|
||||
"UPDATE users SET ref_code=? WHERE tg_id=?", (_generate_ref_code(conn), tg_id)
|
||||
)
|
||||
row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
return dict(row)
|
||||
token = secrets.token_hex(16)
|
||||
ref_code = _generate_ref_code(conn)
|
||||
conn.execute(
|
||||
"INSERT INTO users (tg_id, token, username, created_at) VALUES (?,?,?,?)",
|
||||
(tg_id, token, username, now_iso()),
|
||||
"INSERT INTO users (tg_id, token, username, ref_code, created_at) VALUES (?,?,?,?,?)",
|
||||
(tg_id, token, username, ref_code, now_iso()),
|
||||
)
|
||||
row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
return dict(row)
|
||||
|
||||
|
||||
def get_user_by_ref_code(ref_code: str):
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT * FROM users WHERE ref_code=?", (ref_code,)).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
|
||||
def set_referred_by(tg_id: int, referrer_tg_id: int) -> bool:
|
||||
"""First-touch attribution: only takes effect for a brand-new account
|
||||
(no subscriptions yet) that isn't already attributed, and never to self."""
|
||||
if tg_id == referrer_tg_id:
|
||||
return False
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT referred_by FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
if not row or row["referred_by"] is not None:
|
||||
return False
|
||||
has_sub = conn.execute("SELECT 1 FROM subscriptions WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
if has_sub:
|
||||
return False
|
||||
referrer = conn.execute("SELECT 1 FROM users WHERE tg_id=?", (referrer_tg_id,)).fetchone()
|
||||
if not referrer:
|
||||
return False
|
||||
conn.execute("UPDATE users SET referred_by=? WHERE tg_id=?", (referrer_tg_id, tg_id))
|
||||
return True
|
||||
|
||||
|
||||
def _apply_bonus_days(conn, tg_id: int, days: int):
|
||||
if days <= 0:
|
||||
return
|
||||
row = conn.execute(
|
||||
"SELECT uuid, expires_at FROM subscriptions WHERE tg_id=? AND active=1 AND held_at IS NULL "
|
||||
"ORDER BY expires_at DESC LIMIT 1",
|
||||
(tg_id,),
|
||||
).fetchone()
|
||||
if row:
|
||||
new_expires = datetime.datetime.fromisoformat(row["expires_at"]) + datetime.timedelta(days=days)
|
||||
conn.execute("UPDATE subscriptions SET expires_at=? WHERE uuid=?", (new_expires.isoformat(), row["uuid"]))
|
||||
else:
|
||||
conn.execute(
|
||||
"UPDATE users SET bonus_days_pending = COALESCE(bonus_days_pending, 0) + ? WHERE tg_id=?",
|
||||
(days, tg_id),
|
||||
)
|
||||
|
||||
|
||||
def credit_bonus_days(tg_id: int, days: int):
|
||||
with get_conn() as conn:
|
||||
_apply_bonus_days(conn, tg_id, days)
|
||||
|
||||
|
||||
def referral_stats(tg_id: int) -> dict:
|
||||
with get_conn() as conn:
|
||||
user = conn.execute("SELECT ref_code, bonus_days_pending FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
count = conn.execute(
|
||||
"SELECT COUNT(*) c FROM users WHERE referred_by=? AND referral_rewarded=1", (tg_id,)
|
||||
).fetchone()["c"]
|
||||
return {
|
||||
"ref_code": user["ref_code"] if user else None,
|
||||
"bonus_days_pending": user["bonus_days_pending"] if user else 0,
|
||||
"referred_count": count,
|
||||
}
|
||||
|
||||
|
||||
def get_user_by_token(token: str):
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT * FROM users WHERE token=?", (token,)).fetchone()
|
||||
|
|
@ -377,17 +561,32 @@ def get_user_by_token(token: str):
|
|||
|
||||
def create_subscription(tg_id: int, node: str, plan_days: int, plan_code: str, source: str = "bot", client_uuid: str | None = None):
|
||||
import uuid as uuidlib
|
||||
import config
|
||||
|
||||
cid = client_uuid or str(uuidlib.uuid4())
|
||||
created = datetime.datetime.utcnow()
|
||||
expires = created + datetime.timedelta(days=plan_days)
|
||||
with get_conn() as conn:
|
||||
is_first = conn.execute("SELECT 1 FROM subscriptions WHERE tg_id=?", (tg_id,)).fetchone() is None
|
||||
urow = conn.execute(
|
||||
"SELECT referred_by, referral_rewarded, bonus_days_pending FROM users WHERE tg_id=?", (tg_id,)
|
||||
).fetchone()
|
||||
pending = urow["bonus_days_pending"] if urow else 0
|
||||
if pending:
|
||||
expires += datetime.timedelta(days=pending)
|
||||
conn.execute("UPDATE users SET bonus_days_pending=0 WHERE tg_id=?", (tg_id,))
|
||||
conn.execute(
|
||||
"INSERT INTO subscriptions (uuid, tg_id, node, plan, created_at, expires_at, active, source) "
|
||||
"VALUES (?,?,?,?,?,?,1,?)",
|
||||
(cid, tg_id, node, plan_code, created.isoformat(), expires.isoformat(), source),
|
||||
)
|
||||
return {"uuid": cid, "tg_id": tg_id, "node": node, "plan": plan_code, "expires_at": expires.isoformat()}
|
||||
if is_first and urow and urow["referred_by"] and not urow["referral_rewarded"] and config.REFERRAL_ENABLED:
|
||||
conn.execute("UPDATE users SET referral_rewarded=1 WHERE tg_id=?", (tg_id,))
|
||||
bonus = config.REFERRAL_BONUS_DAYS
|
||||
_apply_bonus_days(conn, tg_id, bonus)
|
||||
_apply_bonus_days(conn, urow["referred_by"], bonus)
|
||||
expires_final = conn.execute("SELECT expires_at FROM subscriptions WHERE uuid=?", (cid,)).fetchone()["expires_at"]
|
||||
return {"uuid": cid, "tg_id": tg_id, "node": node, "plan": plan_code, "expires_at": expires_final}
|
||||
|
||||
|
||||
def list_active_subscriptions(tg_id: int | None = None, node: str | None = None):
|
||||
|
|
@ -645,6 +844,32 @@ def get_subscription(client_uuid: str):
|
|||
return dict(row) if row else None
|
||||
|
||||
|
||||
def list_users(q: str = "", limit: int = 200):
|
||||
limit = max(1, min(int(limit), 1000))
|
||||
q = (q or "").strip()
|
||||
where = ""
|
||||
params = [now_iso(), now_iso()]
|
||||
if q:
|
||||
if q.lstrip("-").isdigit():
|
||||
where = "WHERE u.tg_id = ? OR instr(lower(COALESCE(u.username, '')), ?) > 0"
|
||||
params += [int(q), q.lower()]
|
||||
else:
|
||||
where = "WHERE instr(lower(COALESCE(u.username, '')), ?) > 0"
|
||||
params.append(q.lower().lstrip("@"))
|
||||
params.append(limit)
|
||||
query = (
|
||||
"SELECT u.tg_id, u.username, u.created_at, "
|
||||
"(SELECT COUNT(*) FROM subscriptions s WHERE s.tg_id=u.tg_id) AS subs_total, "
|
||||
"(SELECT COUNT(*) FROM subscriptions s WHERE s.tg_id=u.tg_id AND s.active=1 AND s.held_at IS NULL AND s.expires_at > ?) AS subs_active, "
|
||||
"(SELECT MAX(s.expires_at) FROM subscriptions s WHERE s.tg_id=u.tg_id AND s.active=1 AND s.held_at IS NULL AND s.expires_at > ?) AS active_until, "
|
||||
"(SELECT COUNT(*) FROM devices d WHERE d.tg_id=u.tg_id) AS devices "
|
||||
"FROM users u " + where + " ORDER BY u.created_at DESC LIMIT ?"
|
||||
)
|
||||
with get_conn() as conn:
|
||||
rows = conn.execute(query, params).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def get_user(tg_id: int):
|
||||
with get_conn() as conn:
|
||||
row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone()
|
||||
|
|
@ -676,12 +901,16 @@ def stats():
|
|||
total_subs = conn.execute("SELECT COUNT(*) c FROM subscriptions").fetchone()["c"]
|
||||
gifts_created = conn.execute("SELECT COUNT(*) c FROM gift_codes").fetchone()["c"]
|
||||
gifts_used = conn.execute("SELECT COUNT(*) c FROM gift_codes WHERE used_by IS NOT NULL").fetchone()["c"]
|
||||
nodes_n = conn.execute("SELECT COUNT(*) c FROM nodes WHERE enabled=1").fetchone()["c"]
|
||||
chains_n = conn.execute("SELECT COUNT(*) c FROM chains WHERE enabled=1").fetchone()["c"]
|
||||
return {
|
||||
"users": users_n,
|
||||
"active_subscriptions": active_n,
|
||||
"total_subscriptions": total_subs,
|
||||
"gifts_created": gifts_created,
|
||||
"gifts_used": gifts_used,
|
||||
"nodes": nodes_n,
|
||||
"chains": chains_n,
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -3,7 +3,7 @@ set -e
|
|||
set -o pipefail
|
||||
|
||||
MIRROR_URL="https://api.savsis.xyz/git/mbs-panel.git/"
|
||||
REPO_URL="https://github.com/devsavsis/mbs-panel.git"
|
||||
REPO_URL="https://github.com/savsisbtw/mbs-panel.git"
|
||||
APP_DIR="/opt/mbs-panel"
|
||||
WEBROOT="/var/www/certbot"
|
||||
|
||||
|
|
|
|||
|
|
@ -155,20 +155,20 @@
|
|||
<a href="#features">Возможности</a>
|
||||
<a href="#comparison">Сравнение</a>
|
||||
<a href="#install">Установка</a>
|
||||
<a href="https://github.com/devsavsis/mbs-panel" target="_blank">GitHub</a>
|
||||
<a href="https://github.com/savsisbtw/mbs-panel" target="_blank">GitHub</a>
|
||||
</nav>
|
||||
</header>
|
||||
|
||||
<section class="hero" style="padding-bottom:0">
|
||||
<div class="badges reveal">
|
||||
<img src="https://img.shields.io/github/actions/workflow/status/devsavsis/mbs-panel/ci.yml?label=CI&style=flat-square&color=7c6cf0" alt="CI">
|
||||
<img src="https://img.shields.io/github/license/devsavsis/mbs-panel?style=flat-square&color=7c6cf0" alt="License">
|
||||
<img src="https://img.shields.io/github/stars/devsavsis/mbs-panel?style=flat-square&color=7c6cf0" alt="Stars">
|
||||
<img src="https://img.shields.io/github/actions/workflow/status/savsisbtw/mbs-panel/ci.yml?label=CI&style=flat-square&color=7c6cf0" alt="CI">
|
||||
<img src="https://img.shields.io/github/license/savsisbtw/mbs-panel?style=flat-square&color=7c6cf0" alt="License">
|
||||
<img src="https://img.shields.io/github/stars/savsisbtw/mbs-panel?style=flat-square&color=7c6cf0" alt="Stars">
|
||||
</div>
|
||||
<h1 class="reveal">VPN-панель, которую<br>можно <span class="accent">понять за вечер</span></h1>
|
||||
<p class="reveal">Xray-core, Reality, Hysteria2, Telegram-бот и платежи — в одном небольшом репозитории на Python. Без Docker, без чужой закрытой панели под капотом, без разбора чужого фреймворка перед первым коммитом.</p>
|
||||
<div class="hero-ctas reveal">
|
||||
<a class="btn" href="https://github.com/devsavsis/mbs-panel" target="_blank">Смотреть на GitHub</a>
|
||||
<a class="btn" href="https://github.com/savsisbtw/mbs-panel" target="_blank">Смотреть на GitHub</a>
|
||||
<a class="btn ghost" href="#install">Установка</a>
|
||||
</div>
|
||||
<div class="install reveal">
|
||||
|
|
@ -233,6 +233,11 @@
|
|||
<tr><td>Мультиадминство (раздельные логины)</td><td class="us yes">✓</td><td class="no">—</td><td class="no">в разработке</td></tr>
|
||||
<tr><td>2FA на вход в админку</td><td class="us yes">✓ TOTP</td><td>есть (passkeys/OAuth)</td><td class="no">—</td></tr>
|
||||
<tr><td>Rate-limit на вход/2FA-код</td><td class="us yes">✓</td><td class="no">не документировано</td><td class="no">не документировано</td></tr>
|
||||
<tr><td>Свой путь входа в админку</td><td class="us yes">✓</td><td class="yes">заявлено</td><td class="no">—</td></tr>
|
||||
<tr><td>Цепочки серверов (клиент → A → B → интернет)</td><td class="us yes">✓ мышкой, с замером задержки</td><td>руками в конфиге Xray</td><td>руками в конфиге Xray</td></tr>
|
||||
<tr><td>Сортировка нод мышкой</td><td class="us yes">✓</td><td class="yes">Web UI</td><td class="no">только через конфиг Xray</td></tr>
|
||||
<tr><td>Пауза подписки без потери оплаченных дней</td><td class="us yes">✓</td><td class="no">—</td><td class="yes">есть</td></tr>
|
||||
<tr><td>Исходящие вебхуки (пользователи + ноды)</td><td class="us yes">✓</td><td class="yes">✓</td><td class="no">только пользователи</td></tr>
|
||||
<tr><td>Лицензия</td><td class="us">MIT</td><td>AGPL-3.0</td><td>AGPL-3.0</td></tr>
|
||||
</tbody>
|
||||
</table>
|
||||
|
|
@ -292,16 +297,16 @@
|
|||
<h2>Открытый исходник, MIT</h2>
|
||||
<p class="section-sub">Изначально писалось под конкретный проект — получилось достаточно универсально, чтобы выложить как есть. Issues и PR приветствуются.</p>
|
||||
<div class="cta-row">
|
||||
<a class="btn" href="https://github.com/devsavsis/mbs-panel" target="_blank">github.com/devsavsis/mbs-panel</a>
|
||||
<a class="btn ghost" href="https://github.com/devsavsis/mbs-panel#readme" target="_blank">Читать README</a>
|
||||
<a class="btn" href="https://github.com/savsisbtw/mbs-panel" target="_blank">github.com/savsisbtw/mbs-panel</a>
|
||||
<a class="btn ghost" href="https://github.com/savsisbtw/mbs-panel#readme" target="_blank">Читать README</a>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<footer>
|
||||
<div class="wrap">
|
||||
<div>MBS Panel — made by <a href="https://github.com/devsavsis" target="_blank">savsis</a></div>
|
||||
<div><a href="https://github.com/devsavsis/mbs-panel/blob/main/LICENSE" target="_blank">MIT License</a></div>
|
||||
<div>MBS Panel — made by <a href="https://github.com/savsisbtw" target="_blank">savsis</a></div>
|
||||
<div><a href="https://github.com/savsisbtw/mbs-panel/blob/main/LICENSE" target="_blank">MIT License</a></div>
|
||||
</div>
|
||||
</footer>
|
||||
|
||||
|
|
|
|||
20
links.py
20
links.py
|
|
@ -89,6 +89,17 @@ def vless_uris_for_node(client_uuid: str, node: dict, base_name: str) -> list[st
|
|||
)]
|
||||
|
||||
|
||||
def chain_remark(entry_node: dict, exit_node: dict) -> str:
|
||||
return f"{display_name(entry_node['label'])} → {display_name(exit_node['label'])}"
|
||||
|
||||
|
||||
def chain_uri(client_uuid: str, entry_node: dict, chain: dict, remark: str) -> str:
|
||||
return _tcp_reality_uri(
|
||||
client_uuid, entry_node["address"], chain["port"], entry_node["public_key"],
|
||||
chain["short_id"], entry_node["sni"], "xtls-rprx-vision", remark,
|
||||
)
|
||||
|
||||
|
||||
def build_subscription_text(subs: list[dict]) -> str:
|
||||
import db
|
||||
|
||||
|
|
@ -99,6 +110,10 @@ def build_subscription_text(subs: list[dict]) -> str:
|
|||
best_by_node[s["node"]] = s
|
||||
|
||||
nodes_by_code = {n["code"]: n for n in db.list_nodes()}
|
||||
chains_by_entry = {}
|
||||
for chain in db.list_chains(enabled_only=True):
|
||||
chains_by_entry.setdefault(chain["entry_node"], []).append(chain)
|
||||
|
||||
lines = []
|
||||
for node_code, s in best_by_node.items():
|
||||
node = nodes_by_code.get(node_code)
|
||||
|
|
@ -109,5 +124,10 @@ def build_subscription_text(subs: list[dict]) -> str:
|
|||
hy = hysteria_uri_for_node(node, base_name)
|
||||
if hy:
|
||||
lines.append(hy)
|
||||
for chain in chains_by_entry.get(node_code, []):
|
||||
exit_node = nodes_by_code.get(chain["exit_node"])
|
||||
if not node["enabled"] or not exit_node or not exit_node["enabled"]:
|
||||
continue
|
||||
lines.append(chain_uri(s["uuid"], node, chain, chain_remark(node, exit_node)))
|
||||
raw = "\n".join(lines)
|
||||
return base64.b64encode(raw.encode()).decode()
|
||||
|
|
|
|||
194
mbs
194
mbs
|
|
@ -13,7 +13,10 @@ mbs — управление MBS Panel
|
|||
mbs restart перезапустить всё (bot, api, xray, reload nginx)
|
||||
mbs logs [bot|api|xray] последние строки лога (по умолчанию api)
|
||||
mbs domain показать текущий домен панели
|
||||
mbs update обновить код (зеркало api.savsis.xyz, потом GitHub) и перезапустить (не трогает .env и базу)
|
||||
mbs backup полная копия панели (база, .env, ручные правки, конфиг Xray) в /root/mbs-backups
|
||||
mbs update [ссылка] обновить код и перезапустить (не трогает .env и базу, перед этим сам делает копию). Без ссылки: своё зеркало -> api.savsis.xyz -> GitHub.
|
||||
Со ссылкой на git-репозиторий (зеркало) — берёт обновление оттуда, один раз
|
||||
mbs mirror [ссылка|off] показать / запомнить / убрать своё зеркало, которое mbs update проверяет первым
|
||||
EOF
|
||||
}
|
||||
|
||||
|
|
@ -55,34 +58,172 @@ cmd_domain() {
|
|||
grep "^PANEL_DOMAIN=" "$ENV_FILE"
|
||||
}
|
||||
|
||||
cmd_update() {
|
||||
cd "$APP_DIR"
|
||||
echo "проверяю обновления..."
|
||||
local remote="origin"
|
||||
if ! git fetch --quiet origin main 2>/dev/null; then
|
||||
if git remote | grep -q '^github$'; then
|
||||
echo "зеркало недоступно, пробую github..."
|
||||
if ! git fetch --quiet github main 2>/dev/null; then
|
||||
echo "не удалось получить обновления ни с зеркала, ни с github"
|
||||
MIRROR_FILE="$APP_DIR/.update_mirror"
|
||||
BACKUP_DIR="${MBS_BACKUP_DIR:-/root/mbs-backups}"
|
||||
XRAY_CONFIG="/usr/local/etc/xray/config.json"
|
||||
BACKUP_FILE=""
|
||||
UPDATE_STASHED=0
|
||||
|
||||
snapshot_backup() {
|
||||
local stamp dbsnap file
|
||||
local -a targs=(--exclude=venv --exclude=__pycache__)
|
||||
stamp=$(date +%Y%m%d-%H%M%S)
|
||||
file="$BACKUP_DIR/mbs-$1-$stamp.tar.gz"
|
||||
dbsnap="$APP_DIR/.mbs.db.snapshot"
|
||||
mkdir -p "$BACKUP_DIR" || return 1
|
||||
chmod 700 "$BACKUP_DIR"
|
||||
rm -f "$dbsnap"
|
||||
if [ -f "$APP_DIR/mbs.db" ] && [ -x "$APP_DIR/venv/bin/python" ] && \
|
||||
"$APP_DIR/venv/bin/python" -c "import sqlite3,sys; s=sqlite3.connect(sys.argv[1]); d=sqlite3.connect(sys.argv[2]); s.backup(d); d.close(); s.close()" "$APP_DIR/mbs.db" "$dbsnap" 2>/dev/null; then
|
||||
targs+=(--exclude=mbs.db --exclude=mbs.db-wal --exclude=mbs.db-shm "--transform=s#\.mbs\.db\.snapshot#mbs.db#")
|
||||
fi
|
||||
if [ -f "$XRAY_CONFIG" ]; then
|
||||
tar czf "$file" "${targs[@]}" "$APP_DIR" "$XRAY_CONFIG" 2>/dev/null
|
||||
else
|
||||
tar czf "$file" "${targs[@]}" "$APP_DIR" 2>/dev/null
|
||||
fi
|
||||
local rc=$?
|
||||
rm -f "$dbsnap"
|
||||
if [ $rc -ne 0 ] || [ ! -s "$file" ]; then
|
||||
rm -f "$file"
|
||||
return 1
|
||||
fi
|
||||
chmod 600 "$file"
|
||||
ls -1t "$BACKUP_DIR"/mbs-*.tar.gz 2>/dev/null | tail -n +6 | while read -r old; do rm -f "$old"; done
|
||||
BACKUP_FILE="$file"
|
||||
return 0
|
||||
}
|
||||
|
||||
cmd_backup() {
|
||||
if snapshot_backup manual; then
|
||||
echo "копия сохранена: $BACKUP_FILE"
|
||||
echo "внутри: база (консистентный снапшот), .env, код с твоими правками, конфиг Xray. Хранится 5 последних."
|
||||
else
|
||||
echo "не удалось сделать копию в $BACKUP_DIR (место на диске?)"
|
||||
return 1
|
||||
fi
|
||||
}
|
||||
|
||||
valid_url() {
|
||||
case "$1" in
|
||||
https://*|http://*|ssh://*|git@*:*) ;;
|
||||
*) return 1 ;;
|
||||
esac
|
||||
case "$1" in
|
||||
*[[:space:]]*) return 1 ;;
|
||||
esac
|
||||
return 0
|
||||
}
|
||||
|
||||
fetch_url() {
|
||||
git -c protocol.ext.allow=never fetch --quiet -- "$1" main 2>/dev/null
|
||||
}
|
||||
|
||||
restore_stash() {
|
||||
if [ "$UPDATE_STASHED" = "1" ]; then
|
||||
UPDATE_STASHED=0
|
||||
if git stash pop --quiet 2>/dev/null; then
|
||||
echo "ручные правки вернул на место"
|
||||
else
|
||||
echo "ручные правки остались в git stash (git stash list)"
|
||||
fi
|
||||
fi
|
||||
}
|
||||
|
||||
cmd_mirror() {
|
||||
local url="$1"
|
||||
case "$url" in
|
||||
"")
|
||||
if [ -f "$MIRROR_FILE" ]; then
|
||||
echo "своё зеркало для обновлений: $(head -n 1 "$MIRROR_FILE")"
|
||||
else
|
||||
echo "своё зеркало не задано — mbs update берёт api.savsis.xyz, потом GitHub"
|
||||
fi
|
||||
;;
|
||||
off|clear)
|
||||
rm -f "$MIRROR_FILE"
|
||||
echo "своё зеркало убрано"
|
||||
;;
|
||||
*)
|
||||
if ! valid_url "$url"; then
|
||||
echo "это не похоже на ссылку на git-репозиторий (нужна https://..., ssh://... или git@хост:путь)"
|
||||
return 1
|
||||
fi
|
||||
remote="github"
|
||||
printf '%s\n' "$url" > "$MIRROR_FILE"
|
||||
echo "зеркало запомнил: $url — mbs update теперь проверяет его первым"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
cmd_update() {
|
||||
local arg_url="$1" target="" saved="" before="" stamp="" patch=""
|
||||
cd "$APP_DIR"
|
||||
echo "проверяю обновления..."
|
||||
if [ -n "$arg_url" ]; then
|
||||
if ! valid_url "$arg_url"; then
|
||||
echo "это не похоже на ссылку на git-репозиторий (нужна https://..., ssh://... или git@хост:путь)"
|
||||
return 1
|
||||
fi
|
||||
if ! fetch_url "$arg_url"; then
|
||||
echo "не удалось получить обновления по ссылке: $arg_url"
|
||||
return 1
|
||||
fi
|
||||
target=$(git rev-parse FETCH_HEAD)
|
||||
echo "источник: $arg_url"
|
||||
else
|
||||
if [ -f "$MIRROR_FILE" ]; then
|
||||
saved=$(head -n 1 "$MIRROR_FILE" | tr -d '[:space:]')
|
||||
fi
|
||||
if [ -n "$saved" ] && valid_url "$saved" && fetch_url "$saved"; then
|
||||
target=$(git rev-parse FETCH_HEAD)
|
||||
echo "источник: своё зеркало $saved"
|
||||
elif git fetch --quiet origin main 2>/dev/null; then
|
||||
target=$(git rev-parse origin/main)
|
||||
elif git remote | grep -q '^github$' && git fetch --quiet github main 2>/dev/null; then
|
||||
echo "зеркало недоступно, взял с github..."
|
||||
target=$(git rev-parse github/main)
|
||||
else
|
||||
echo "не удалось получить обновления"
|
||||
echo "не удалось получить обновления ни с зеркала, ни с github"
|
||||
return 1
|
||||
fi
|
||||
fi
|
||||
local before after
|
||||
before=$(git rev-parse HEAD)
|
||||
after=$(git rev-parse "$remote/main")
|
||||
if [ "$before" = "$after" ]; then
|
||||
if [ "$before" = "$target" ]; then
|
||||
echo "уже последняя версия ($before)."
|
||||
return 0
|
||||
fi
|
||||
if git merge-base --is-ancestor "$target" "$before" 2>/dev/null; then
|
||||
echo "на сервере версия новее, чем в источнике ($before) — ничего не делаю."
|
||||
return 0
|
||||
fi
|
||||
echo "текущая: $before"
|
||||
echo "новая: $after"
|
||||
if ! git merge --ff-only "$remote/main" --quiet; then
|
||||
echo "не вышло быстро обновиться (похоже, файлы правились вручную на сервере) — разберись руками: git status"
|
||||
echo "новая: $target"
|
||||
|
||||
if ! snapshot_backup before-update; then
|
||||
echo "не вышло сделать резервную копию в $BACKUP_DIR — обновление не начинаю, чтобы ничего не потерять (место на диске?)"
|
||||
return 1
|
||||
fi
|
||||
echo "резервная копия перед обновлением: $BACKUP_FILE"
|
||||
|
||||
if [ -n "$(git status --porcelain --untracked-files=no)" ]; then
|
||||
stamp=$(date +%Y%m%d-%H%M%S)
|
||||
mkdir -p "$APP_DIR/local-changes"
|
||||
patch="$APP_DIR/local-changes/local-changes-$stamp.patch"
|
||||
git diff HEAD > "$patch"
|
||||
if GIT_AUTHOR_NAME=mbs GIT_AUTHOR_EMAIL=mbs@localhost GIT_COMMITTER_NAME=mbs GIT_COMMITTER_EMAIL=mbs@localhost git stash push --quiet -m "mbs-update-$stamp"; then
|
||||
UPDATE_STASHED=1
|
||||
echo "на сервере были ручные правки — убрал в сторону, ничего не потеряно:"
|
||||
echo " патч: $patch"
|
||||
echo " stash: git stash list (вернуть обратно: git stash pop)"
|
||||
else
|
||||
echo "не вышло спрятать ручные правки, остановился — разберись руками: git status"
|
||||
return 1
|
||||
fi
|
||||
fi
|
||||
|
||||
if ! git merge --ff-only "$target" --quiet 2>/dev/null; then
|
||||
echo "не вышло быстро обновиться (на сервере есть свои коммиты, которых нет в источнике) — разберись руками: git log"
|
||||
restore_stash
|
||||
return 1
|
||||
fi
|
||||
echo "обновляю зависимости..."
|
||||
|
|
@ -92,6 +233,7 @@ cmd_update() {
|
|||
echo "новый код не проходит проверку, откатываюсь на $before..."
|
||||
git reset --hard "$before" --quiet
|
||||
venv/bin/pip install --quiet -r requirements.txt
|
||||
restore_stash
|
||||
return 1
|
||||
fi
|
||||
echo "обновляю сам CLI..."
|
||||
|
|
@ -115,22 +257,36 @@ cmd_update() {
|
|||
sleep 2
|
||||
if systemctl is-active --quiet mbs-bot && systemctl is-active --quiet mbs-api; then
|
||||
echo "обновлено: $before -> $(git rev-parse --short HEAD)"
|
||||
if [ "$UPDATE_STASHED" = "1" ]; then
|
||||
echo "твои ручные правки лежат в git stash и в $patch — если они нужны, посмотри git stash show -p"
|
||||
fi
|
||||
echo "если что-то пошло не так: копия $BACKUP_FILE (распаковать: tar xzf файл -C /)"
|
||||
if [ -n "$arg_url" ]; then
|
||||
echo "чтобы всегда обновляться с этого зеркала: mbs mirror $arg_url"
|
||||
fi
|
||||
else
|
||||
echo "сервисы не поднялись после обновления, откатываюсь на $before..."
|
||||
git reset --hard "$before" --quiet
|
||||
venv/bin/pip install --quiet -r requirements.txt
|
||||
systemctl restart mbs-bot mbs-api
|
||||
restore_stash
|
||||
echo "откачено обратно на $before"
|
||||
return 1
|
||||
fi
|
||||
}
|
||||
|
||||
main() {
|
||||
case "$1" in
|
||||
pass) cmd_pass "$2" ;;
|
||||
status) cmd_status ;;
|
||||
restart) cmd_restart ;;
|
||||
logs) cmd_logs "$2" ;;
|
||||
domain) cmd_domain ;;
|
||||
update) cmd_update ;;
|
||||
update) cmd_update "$2" ;;
|
||||
mirror) cmd_mirror "$2" ;;
|
||||
backup) cmd_backup ;;
|
||||
*) usage ;;
|
||||
esac
|
||||
}
|
||||
|
||||
main "$@"; exit $?
|
||||
|
|
|
|||
179
nodeprov.py
179
nodeprov.py
|
|
@ -1,3 +1,6 @@
|
|||
import contextlib
|
||||
import fcntl
|
||||
import hashlib
|
||||
import json
|
||||
import secrets
|
||||
import socket
|
||||
|
|
@ -5,12 +8,12 @@ import subprocess
|
|||
|
||||
import paramiko
|
||||
|
||||
import chains
|
||||
from config import PANEL_DOMAIN
|
||||
|
||||
MGMT_KEY_PATH = "/root/.ssh/mbs_nodes_ed25519"
|
||||
MGMT_KNOWN_HOSTS_PATH = "/root/.ssh/mbs_nodes_known_hosts"
|
||||
LOCAL_TAGS = {"vless-tcp-reality", "vless-grpc-reality", "vless-xhttp-reality", "vless-ws-tls"}
|
||||
TAG_FLOW = {"vless-tcp-reality": "xtls-rprx-vision"}
|
||||
REMOTE_CONFIG_PATH = "/usr/local/etc/xray/config.json"
|
||||
|
||||
ONE_COMMAND_TEMPLATE = "bash <(curl -Ls https://{panel}/install/{token}.sh)"
|
||||
|
||||
|
|
@ -68,6 +71,28 @@ set -e
|
|||
echo "== MBS Panel node install =="
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
|
||||
PREFLIGHT_FAIL=0
|
||||
if ! {{ [ -f /usr/local/etc/xray/config.json ] && grep -q mbs-grpc /usr/local/etc/xray/config.json; }}; then
|
||||
if [ -d /usr/local/bin/xray ]; then
|
||||
echo "СТОП: /usr/local/bin/xray это каталог, на сервере уже стоит чужой прокси (Marzban-node и подобное)"
|
||||
PREFLIGHT_FAIL=1
|
||||
fi
|
||||
if command -v docker >/dev/null 2>&1 && docker ps --format '{{{{.Names}}}}' 2>/dev/null | grep -qiE 'marzban|xray|remnawave|v2ray|hysteria'; then
|
||||
echo "СТОП: в Docker на этом сервере уже крутится прокси"
|
||||
PREFLIGHT_FAIL=1
|
||||
fi
|
||||
for p in {ports}; do
|
||||
if ss -ltn 2>/dev/null | awk '{{print $4}}' | grep -qE "[:.]$p$"; then
|
||||
echo "СТОП: порт $p уже занят другим процессом"
|
||||
PREFLIGHT_FAIL=1
|
||||
fi
|
||||
done
|
||||
fi
|
||||
if [ "$PREFLIGHT_FAIL" = "1" ]; then
|
||||
echo "Нужен чистый сервер. Ничего не установлено и не изменено, ключ панели не добавлен."
|
||||
exit 1
|
||||
fi
|
||||
|
||||
mkdir -p /root/.ssh
|
||||
chmod 700 /root/.ssh
|
||||
curl -Ls https://{panel}/mgmt-pubkey.txt >> /root/.ssh/authorized_keys
|
||||
|
|
@ -88,9 +113,19 @@ XRAYCFG
|
|||
command -v ufw >/dev/null 2>&1 && {{ ufw allow 22/tcp || true; {ufw_rules} }}
|
||||
systemctl enable xray >/dev/null 2>&1 || true
|
||||
systemctl restart xray
|
||||
sleep 1
|
||||
STATUS=$(systemctl is-active xray)
|
||||
OK=0
|
||||
for i in 1 2 3 4 5 6 7 8; do
|
||||
sleep 1
|
||||
if systemctl is-active --quiet xray && ss -ltn 2>/dev/null | awk '{{print $4}}' | grep -qE "[:.]{first_port}$"; then
|
||||
OK=$((OK+1))
|
||||
else
|
||||
OK=0
|
||||
fi
|
||||
if [ "$OK" -ge 3 ]; then break; fi
|
||||
done
|
||||
if [ "$OK" -ge 3 ]; then STATUS=active; else STATUS=failed; fi
|
||||
echo "xray status: $STATUS"
|
||||
if [ "$STATUS" != "active" ]; then journalctl -u xray -n 15 --no-pager 2>/dev/null || true; fi
|
||||
|
||||
{hysteria_block}
|
||||
MY_IP=$(curl -s https://api.ipify.org || echo unknown)
|
||||
|
|
@ -231,9 +266,11 @@ def render_install_script(node: dict) -> str:
|
|||
sni=node["sni"],
|
||||
)
|
||||
|
||||
ports = " ".join(str(t["port"]) for t in transports)
|
||||
return SELF_INSTALL_SCRIPT.format(
|
||||
panel=PANEL_DOMAIN, token=node["provision_token"], config_json=config_json,
|
||||
certbot_block=certbot_block, hysteria_block=hysteria_block, ufw_rules=ufw_rules,
|
||||
ports=ports, first_port=transports[0]["port"],
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -254,26 +291,45 @@ class RemoteConfigError(Exception):
|
|||
pass
|
||||
|
||||
|
||||
def _remote_edit_clients(node: dict, mutate_fn):
|
||||
def _busy_ports(client) -> set:
|
||||
_, stdout, _ = client.exec_command("ss -ltnH 2>/dev/null | awk '{print $4}'", timeout=10)
|
||||
busy = set()
|
||||
for line in stdout.read().decode(errors="replace").splitlines():
|
||||
tail = line.rsplit(":", 1)[-1]
|
||||
if tail.isdigit():
|
||||
busy.add(int(tail))
|
||||
return busy
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _node_lock(address: str):
|
||||
key = hashlib.sha1(address.encode()).hexdigest()[:12]
|
||||
with open(f"/tmp/mbs-node-{key}.lock", "w") as lock_file:
|
||||
fcntl.flock(lock_file, fcntl.LOCK_EX)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
fcntl.flock(lock_file, fcntl.LOCK_UN)
|
||||
|
||||
|
||||
def _remote_edit_config(node: dict, mutate_fn):
|
||||
with _node_lock(node["address"]):
|
||||
return _remote_edit_config_locked(node, mutate_fn)
|
||||
|
||||
|
||||
def _remote_edit_config_locked(node: dict, mutate_fn):
|
||||
client = _mgmt_connect(node["address"])
|
||||
try:
|
||||
sftp = client.open_sftp()
|
||||
with sftp.open("/usr/local/etc/xray/config.json") as f:
|
||||
with sftp.open(REMOTE_CONFIG_PATH) as f:
|
||||
cfg = json.loads(f.read().decode())
|
||||
changed = False
|
||||
for ib in cfg["inbounds"]:
|
||||
if ib.get("tag") not in LOCAL_TAGS:
|
||||
continue
|
||||
clients = ib["settings"]["clients"]
|
||||
new_clients = mutate_fn(clients, ib["tag"])
|
||||
if new_clients is not None:
|
||||
ib["settings"]["clients"] = new_clients
|
||||
changed = True
|
||||
if not changed:
|
||||
result = mutate_fn(cfg, client)
|
||||
if not result["changed"]:
|
||||
sftp.close()
|
||||
return
|
||||
return result
|
||||
data = json.dumps(cfg, indent=2).encode()
|
||||
tmp_path = "/usr/local/etc/xray/config.json.validate.tmp"
|
||||
tmp_path = REMOTE_CONFIG_PATH + ".validate.tmp"
|
||||
prev_path = REMOTE_CONFIG_PATH + ".mbs-prev"
|
||||
with sftp.open(tmp_path, "wb") as f:
|
||||
f.write(data)
|
||||
_, stdout, stderr = client.exec_command(f"/usr/local/bin/xray run -test -format=json -config {tmp_path}", timeout=15)
|
||||
|
|
@ -283,23 +339,43 @@ def _remote_edit_clients(node: dict, mutate_fn):
|
|||
client.exec_command(f"rm -f {tmp_path}")
|
||||
sftp.close()
|
||||
raise RemoteConfigError(f"config test failed on {node['address']}: {test_out}")
|
||||
client.exec_command(f"mv {tmp_path} /usr/local/etc/xray/config.json")[1].channel.recv_exit_status()
|
||||
client.exec_command(f"cp -p {REMOTE_CONFIG_PATH} {prev_path}")[1].channel.recv_exit_status()
|
||||
client.exec_command(f"mv {tmp_path} {REMOTE_CONFIG_PATH}")[1].channel.recv_exit_status()
|
||||
sftp.close()
|
||||
_, stdout, stderr = client.exec_command("systemctl restart xray", timeout=20)
|
||||
_, stdout, stderr = client.exec_command("systemctl restart xray && sleep 1 && systemctl is-active xray", timeout=30)
|
||||
restart_exit = stdout.channel.recv_exit_status()
|
||||
if restart_exit != 0:
|
||||
err = stderr.read().decode(errors="replace").strip()
|
||||
raise RemoteConfigError(f"xray restart failed on {node['address']}: {err}")
|
||||
client.exec_command(f"cp -p {prev_path} {REMOTE_CONFIG_PATH} && systemctl restart xray")[1].channel.recv_exit_status()
|
||||
raise RemoteConfigError(f"xray не поднялся на {node['address']}, конфиг откатили назад: {err}")
|
||||
for port in result.get("new_ports") or []:
|
||||
client.exec_command(f"command -v ufw >/dev/null 2>&1 && ufw allow {int(port)}/tcp || true")[1].channel.recv_exit_status()
|
||||
return result
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
|
||||
def _remote_edit_clients(node: dict, mutate_fn):
|
||||
def mutate(cfg, client):
|
||||
changed = False
|
||||
for ib in cfg["inbounds"]:
|
||||
tag = ib.get("tag")
|
||||
if not chains.is_user_tag(tag):
|
||||
continue
|
||||
new_clients = mutate_fn(ib["settings"]["clients"], tag)
|
||||
if new_clients is not None:
|
||||
ib["settings"]["clients"] = new_clients
|
||||
changed = True
|
||||
return {"changed": changed}
|
||||
_remote_edit_config(node, mutate)
|
||||
|
||||
|
||||
def remote_add_client(node: dict, client_uuid: str, email: str):
|
||||
def mutate(clients, tag):
|
||||
if any(c["id"] == client_uuid for c in clients):
|
||||
return None
|
||||
entry = {"id": client_uuid, "email": email}
|
||||
flow = TAG_FLOW.get(tag)
|
||||
flow = chains.flow_for_tag(tag)
|
||||
if flow:
|
||||
entry["flow"] = flow
|
||||
clients.append(entry)
|
||||
|
|
@ -314,24 +390,49 @@ def remote_remove_client(node: dict, client_uuid: str):
|
|||
_remote_edit_clients(node, mutate)
|
||||
|
||||
|
||||
def remote_sync(node: dict, active_subs: list[dict]):
|
||||
active_by_id = {s["uuid"]: s for s in active_subs}
|
||||
def remote_reconcile(node: dict, wanted: dict, relay_wanted: dict, entry_chains: list, exit_nodes: dict, apply_chains: bool = True):
|
||||
def mutate(cfg, client):
|
||||
usable = entry_chains
|
||||
skipped = []
|
||||
if apply_chains:
|
||||
usable, skipped = chains.split_busy_chains(cfg, entry_chains, _busy_ports(client))
|
||||
new_ports = chains.new_ports_needed(cfg, usable) if apply_chains else []
|
||||
changed, problems = chains.sync_config(cfg, wanted, relay_wanted, usable, exit_nodes, apply_chains=apply_chains)
|
||||
return {"changed": changed, "new_ports": new_ports, "problems": skipped + problems}
|
||||
return _remote_edit_config(node, mutate)
|
||||
|
||||
def mutate(clients, tag):
|
||||
current_ids = {c["id"] for c in clients}
|
||||
if current_ids == set(active_by_id.keys()):
|
||||
return None
|
||||
new_clients = [c for c in clients if c["id"] in active_by_id]
|
||||
existing_ids = {c["id"] for c in new_clients}
|
||||
flow = TAG_FLOW.get(tag)
|
||||
for cid in active_by_id:
|
||||
if cid not in existing_ids:
|
||||
entry = {"id": cid, "email": cid}
|
||||
if flow:
|
||||
entry["flow"] = flow
|
||||
new_clients.append(entry)
|
||||
return new_clients
|
||||
_remote_edit_clients(node, mutate)
|
||||
|
||||
PROBE_SCRIPT = """for i in 1 2 3; do
|
||||
s=$(date +%s%N)
|
||||
if timeout 3 bash -c 'exec 3<>/dev/tcp/{host}/{port}' 2>/dev/null; then
|
||||
e=$(date +%s%N)
|
||||
echo $(( (e - s) / 1000000 ))
|
||||
else
|
||||
echo -1
|
||||
fi
|
||||
done
|
||||
"""
|
||||
|
||||
|
||||
def remote_probe(node: dict, host: str, port: int) -> list:
|
||||
if not chains.HOST_RE.match(host or ""):
|
||||
raise ValueError("bad host")
|
||||
port = int(port)
|
||||
client = _mgmt_connect(node["address"])
|
||||
try:
|
||||
stdin, stdout, _ = client.exec_command("bash -s", timeout=30)
|
||||
stdin.write(PROBE_SCRIPT.format(host=host, port=port))
|
||||
stdin.channel.shutdown_write()
|
||||
out = stdout.read().decode(errors="replace")
|
||||
finally:
|
||||
client.close()
|
||||
samples = []
|
||||
for line in out.split():
|
||||
try:
|
||||
samples.append(int(line))
|
||||
except ValueError:
|
||||
continue
|
||||
return samples
|
||||
|
||||
|
||||
def remote_query_stats(node: dict) -> dict:
|
||||
|
|
|
|||
14
settings.py
14
settings.py
|
|
@ -87,3 +87,17 @@ def get_hwid_settings() -> dict:
|
|||
"enabled": _bool(raw.get("HWID_LIMIT_ENABLED"), config.HWID_LIMIT_ENABLED),
|
||||
"fallback_limit": limit if limit > 0 else config.HWID_FALLBACK_LIMIT,
|
||||
}
|
||||
|
||||
|
||||
def get_referral_settings() -> dict:
|
||||
raw = legal.read_env_vars(["REFERRAL_ENABLED", "REFERRAL_BONUS_DAYS"])
|
||||
days = _positive_int(raw.get("REFERRAL_BONUS_DAYS"), config.REFERRAL_BONUS_DAYS)
|
||||
return {
|
||||
"enabled": _bool(raw.get("REFERRAL_ENABLED"), config.REFERRAL_ENABLED),
|
||||
"bonus_days": days if days > 0 else config.REFERRAL_BONUS_DAYS,
|
||||
}
|
||||
|
||||
|
||||
def set_referral_settings(enabled: bool, bonus_days: int):
|
||||
legal.update_env_var("REFERRAL_ENABLED", "true" if enabled else "false")
|
||||
legal.update_env_var("REFERRAL_BONUS_DAYS", str(int(bonus_days)))
|
||||
|
|
|
|||
155
tests/test_mbs_update.sh
Normal file
155
tests/test_mbs_update.sh
Normal file
|
|
@ -0,0 +1,155 @@
|
|||
#!/bin/bash
|
||||
REPO="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
T=$(mktemp -d)
|
||||
cd "$T"
|
||||
export GIT_CONFIG_COUNT=1 GIT_CONFIG_KEY_0=core.autocrlf GIT_CONFIG_VALUE_0=false
|
||||
export GIT_AUTHOR_NAME=t GIT_AUTHOR_EMAIL=t@t GIT_COMMITTER_NAME=t GIT_COMMITTER_EMAIL=t@t
|
||||
mkdir -p stub bin sysd
|
||||
cat > stub/systemctl << 'EOF'
|
||||
#!/bin/sh
|
||||
exit 0
|
||||
EOF
|
||||
chmod +x stub/systemctl
|
||||
export PATH="$T/stub:$PATH"
|
||||
export MBS_BACKUP_DIR="$T/backups"
|
||||
|
||||
PASS=0
|
||||
FAIL=0
|
||||
ok() { echo "PASS $1"; PASS=$((PASS+1)); }
|
||||
bad() { echo "FAIL $1 :: $2"; FAIL=$((FAIL+1)); }
|
||||
check() { if eval "$2"; then ok "$1"; else bad "$1" "$3"; fi; }
|
||||
|
||||
sed -e "s#APP_DIR=\"/opt/mbs-panel\"#APP_DIR=\"$T/app\"#" -e "s#/usr/local/bin/mbs#$T/bin/mbs#g" -e "s#/etc/systemd/system#$T/sysd#g" "$REPO/mbs" > "$T/mbs-run"
|
||||
|
||||
git init -q --bare -b main origin.git
|
||||
git clone -q origin.git seed 2>/dev/null
|
||||
cd seed
|
||||
git checkout -q -b main 2>/dev/null || true
|
||||
mkdir -p systemd
|
||||
cp "$REPO/mbs" mbs
|
||||
echo "v1" > bot.py
|
||||
echo "v1" > config.py
|
||||
echo "v1" > db.py
|
||||
echo "v1" > settings.py
|
||||
echo "x" > requirements.txt
|
||||
echo "[Service]" > systemd/mbs-bot.service
|
||||
echo "ExecStart=x --workers __WORKERS__" > systemd/mbs-api.service
|
||||
cp "$REPO/.gitignore" .gitignore
|
||||
git add -A && git commit -q -m A && git push -q origin main
|
||||
cd "$T"
|
||||
git clone -q origin.git app
|
||||
mkdir -p app/venv/bin
|
||||
printf '#!/bin/sh\nexit 0\n' > app/venv/bin/pip
|
||||
printf '#!/bin/sh\nexec python "$@"\n' > app/venv/bin/python
|
||||
chmod +x app/venv/bin/pip app/venv/bin/python
|
||||
echo "SECRET=1" > app/.env
|
||||
python -c "import sqlite3,sys; c=sqlite3.connect(sys.argv[1]); c.execute('pragma journal_mode=wal'); c.execute('create table t(x)'); c.execute('insert into t values (42)'); c.commit(); c.close()" app/mbs.db
|
||||
|
||||
cd seed && echo "v2" > bot.py && echo "v2" > db.py && git commit -qam B && git push -q origin main && cd "$T"
|
||||
|
||||
echo "manual edit" >> app/bot.py
|
||||
echo "manual edit" >> app/config.py
|
||||
echo "manual edit" >> app/db.py
|
||||
echo "manual edit" >> app/settings.py
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
echo "$OUT" > out1.txt
|
||||
check "update succeeds despite manual edits on the server (the reported bug)" "[ $RC -eq 0 ]" "rc=$RC $OUT"
|
||||
check "bot.py is the new version" "[ \"\$(cat app/bot.py)\" = v2 ]" "$(cat app/bot.py)"
|
||||
check "a stash with the manual edits exists" "[ \$(git -C app stash list | wc -l) -eq 1 ]"
|
||||
check "patch with the manual edits saved" "ls app/local-changes/*.patch >/dev/null 2>&1 && grep -q 'manual edit' app/local-changes/*.patch"
|
||||
check "user is told where the edits went" "grep -q 'ручные правки' out1.txt && grep -q 'патч' out1.txt"
|
||||
check "working tree clean after update" "[ -z \"\$(git -C app status --porcelain --untracked-files=no)\" ]"
|
||||
check "cli copied" "[ -f bin/mbs ]"
|
||||
check "a pre-update backup was made" "[ \$(ls backups/mbs-before-update-*.tar.gz 2>/dev/null | wc -l) -eq 1 ]" "$(ls backups 2>&1)"
|
||||
check "update output points at the backup" "grep -q 'резервная копия перед обновлением' out1.txt"
|
||||
mkdir -p unpack && tar xzf backups/mbs-before-update-*.tar.gz -C unpack
|
||||
APPREL="${T#/}/app"
|
||||
check "backup keeps .env" "grep -q SECRET=1 unpack/$APPREL/.env"
|
||||
check "backup keeps the manual edits as they were before the update" "grep -q 'manual edit' unpack/$APPREL/bot.py && grep -q 'manual edit' unpack/$APPREL/config.py"
|
||||
check "backup database is a consistent sqlite copy" "[ \"\$(python -c \"import sqlite3,sys; print(sqlite3.connect(sys.argv[1]).execute('select x from t').fetchone()[0])\" unpack/$APPREL/mbs.db)\" = 42 ]"
|
||||
check "backup does not drag the venv along" "[ ! -d unpack/$APPREL/venv ]"
|
||||
check "no snapshot temp file is left behind" "[ ! -e app/.mbs.db.snapshot ]"
|
||||
OUT=$(bash "$T/mbs-run" backup 2>&1); RC=$?
|
||||
check "mbs backup makes a manual copy" "[ $RC -eq 0 ] && ls backups/mbs-manual-*.tar.gz >/dev/null 2>&1" "$OUT"
|
||||
for i in 1 2 3 4 5 6 7; do sleep 1.1; bash "$T/mbs-run" backup >/dev/null 2>&1; done
|
||||
check "only the 5 newest backups are kept" "[ \$(ls backups/mbs-*.tar.gz | wc -l) -eq 5 ]" "$(ls backups | wc -l)"
|
||||
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); echo "$OUT" > out2.txt
|
||||
check "second run says already latest" "grep -q 'уже последняя' out2.txt" "$OUT"
|
||||
|
||||
cd "$T/app" && git stash drop -q && git reset -q --hard HEAD; cd "$T"
|
||||
|
||||
git clone -q --bare origin.git mirror2.git
|
||||
cd seed && git pull -q origin main 2>/dev/null; echo "v3" > bot.py && git commit -qam C && git push -q "$T/mirror2.git" main && cd "$T"
|
||||
git -C mirror2.git update-server-info
|
||||
python -m http.server 8799 --directory "$T" > http.log 2>&1 &
|
||||
HTTP_PID=$!
|
||||
sleep 1.5
|
||||
OUT=$(bash "$T/mbs-run" update "http://127.0.0.1:8799/mirror2.git" 2>&1); RC=$?
|
||||
echo "$OUT" > out3.txt
|
||||
check "update by mirror url works" "[ $RC -eq 0 ] && [ \"\$(cat app/bot.py)\" = v3 ]" "rc=$RC $OUT"
|
||||
check "url run mentions the source and how to save it" "grep -q 'источник: http://127.0.0.1:8799/mirror2.git' out3.txt && grep -q 'mbs mirror http' out3.txt" "$OUT"
|
||||
check "url alone is not saved automatically" "[ ! -f app/.update_mirror ]"
|
||||
|
||||
bash "$T/mbs-run" mirror "http://127.0.0.1:8799/mirror2.git" > out4.txt 2>&1
|
||||
check "mirror command saves the url" "grep -q 'http://127.0.0.1:8799/mirror2.git' app/.update_mirror"
|
||||
check "mirror file is git-ignored" "[ -z \"\$(git -C app status --porcelain)\" ]" "$(git -C app status --porcelain)"
|
||||
bash "$T/mbs-run" mirror > out5.txt 2>&1
|
||||
check "mirror shows the saved url" "grep -q 'своё зеркало для обновлений: http://127.0.0.1:8799' out5.txt"
|
||||
|
||||
cd seed && echo "v4" > bot.py && git commit -qam D && git push -q "$T/mirror2.git" main && cd "$T"
|
||||
git -C mirror2.git update-server-info
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
echo "$OUT" > out6.txt
|
||||
check "plain update prefers the saved mirror" "[ $RC -eq 0 ] && grep -q 'источник: своё зеркало' out6.txt && [ \"\$(cat app/bot.py)\" = v4 ]" "rc=$RC $OUT"
|
||||
|
||||
bash "$T/mbs-run" mirror off > out7.txt 2>&1
|
||||
check "mirror off removes it" "[ ! -f app/.update_mirror ]"
|
||||
|
||||
for bad_url in "ext::sh -c 'touch $T/pwned'" "-uHEAD" "file:///etc" "ftp://x/y" "https://a b/c" "--upload-pack=touch $T/pwned2"; do
|
||||
OUT=$(bash "$T/mbs-run" update "$bad_url" 2>&1); RC=$?
|
||||
check "rejects unsafe source '$bad_url'" "[ $RC -ne 0 ] && grep -q 'не похоже на ссылку' <<< \"\$OUT\"" "rc=$RC $OUT"
|
||||
done
|
||||
check "no command was executed via a crafted url" "[ ! -e pwned ] && [ ! -e pwned2 ]"
|
||||
OUT=$(bash "$T/mbs-run" mirror "file:///etc" 2>&1); RC=$?
|
||||
check "mirror refuses unsafe urls too" "[ $RC -ne 0 ] && [ ! -f app/.update_mirror ]"
|
||||
|
||||
OUT=$(bash "$T/mbs-run" update "http://127.0.0.1:1/nope.git" 2>&1); RC=$?
|
||||
check "unreachable mirror fails cleanly" "[ $RC -ne 0 ] && grep -q 'не удалось получить обновления по ссылке' <<< \"\$OUT\"" "rc=$RC $OUT"
|
||||
|
||||
cd seed && echo "v5" > bot.py && echo "broken(" > broken.py && git add -A && git commit -qm E && git push -q origin main && cd "$T"
|
||||
git -C app fetch -q origin main
|
||||
git -C app merge -q --ff-only origin/main~1 2>/dev/null || true
|
||||
cd app && git reset -q --hard origin/main~1 2>/dev/null; cd "$T"
|
||||
echo "local tweak" >> app/settings.py
|
||||
BEFORE=$(git -C app rev-parse HEAD)
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
echo "$OUT" > out8.txt
|
||||
check "broken new code is rolled back" "[ $RC -ne 0 ] && grep -q 'откатываюсь' out8.txt && [ \"\$(git -C app rev-parse HEAD)\" = \"$BEFORE\" ]" "rc=$RC $OUT"
|
||||
check "manual edits are restored after the rollback" "grep -q 'local tweak' app/settings.py && [ \$(git -C app stash list | wc -l) -eq 0 ]" "$(git -C app stash list)"
|
||||
|
||||
cd app && git checkout -q -- . && git reset -q --hard origin/main~1 && cd "$T"
|
||||
cd app && echo "own" > own.txt && git add own.txt && git commit -qm "local commit" && cd "$T"
|
||||
cd seed && git rm -q broken.py && echo "v6" > bot.py && git commit -qam F && git push -q origin main && cd "$T"
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
echo "$OUT" > out9.txt
|
||||
check "diverged server history is refused, not merged" "[ $RC -ne 0 ] && grep -q 'свои коммиты' out9.txt" "rc=$RC $OUT"
|
||||
check "refused update leaves the local commit alone" "git -C app log --oneline | grep -q 'local commit'"
|
||||
|
||||
cd app && git reset -q --hard origin/main && echo "ahead" > ahead.txt && git add ahead.txt && git commit -qm ahead && cd "$T"
|
||||
OUT=$(bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
check "server newer than the source is left alone" "[ $RC -eq 0 ] && grep -q 'версия новее' <<< \"\$OUT\"" "rc=$RC $OUT"
|
||||
|
||||
cd "$T/app" && git reset -q --hard origin/main && cd "$T"
|
||||
cd seed && echo "v7" > bot.py && git commit -qam G && git push -q origin main && cd "$T"
|
||||
BEFORE=$(git -C app rev-parse HEAD)
|
||||
echo "x" > notadir
|
||||
OUT=$(MBS_BACKUP_DIR="$T/notadir/x" bash "$T/mbs-run" update 2>&1); RC=$?
|
||||
check "update refuses to start when no backup can be made" "[ $RC -ne 0 ] && grep -q 'не начинаю' <<< \"\$OUT\" && [ \"\$(git -C app rev-parse HEAD)\" = \"$BEFORE\" ]" "rc=$RC $OUT"
|
||||
check "refused update leaves the code untouched" "[ \"\$(cat app/bot.py)\" != v7 ]"
|
||||
|
||||
kill $HTTP_PID 2>/dev/null
|
||||
cd /
|
||||
rm -rf "$T"
|
||||
echo "RESULT pass=$PASS fail=$FAIL"
|
||||
[ "$FAIL" -eq 0 ]
|
||||
177
xray_manager.py
177
xray_manager.py
|
|
@ -1,14 +1,15 @@
|
|||
import json
|
||||
import os
|
||||
import socket
|
||||
import subprocess
|
||||
import fcntl
|
||||
import contextlib
|
||||
import time
|
||||
|
||||
from config import XRAY_CONFIG_PATH, DE1_TRANSPORTS
|
||||
import chains
|
||||
from config import XRAY_CONFIG_PATH
|
||||
|
||||
_LOCK_PATH = XRAY_CONFIG_PATH + ".lock"
|
||||
_LOCAL_TAGS = {t["tag"] for t in DE1_TRANSPORTS}
|
||||
_TAG_FLOW = {t["tag"]: t.get("flow") for t in DE1_TRANSPORTS}
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
|
|
@ -101,7 +102,7 @@ def _reload_xray():
|
|||
|
||||
|
||||
def _local_inbounds(cfg):
|
||||
return [ib for ib in cfg["inbounds"] if ib.get("tag") in _LOCAL_TAGS]
|
||||
return [ib for ib in cfg["inbounds"] if chains.is_user_tag(ib.get("tag"))]
|
||||
|
||||
|
||||
def add_client(client_uuid: str, email: str):
|
||||
|
|
@ -113,7 +114,7 @@ def add_client(client_uuid: str, email: str):
|
|||
if any(c["id"] == client_uuid for c in clients):
|
||||
continue
|
||||
entry = {"id": client_uuid, "email": email}
|
||||
flow = _TAG_FLOW.get(ib["tag"])
|
||||
flow = chains.flow_for_tag(ib["tag"])
|
||||
if flow:
|
||||
entry["flow"] = flow
|
||||
clients.append(entry)
|
||||
|
|
@ -138,38 +139,138 @@ def remove_client(client_uuid: str):
|
|||
_reload_xray()
|
||||
|
||||
|
||||
def _node_usable(node):
|
||||
return bool(node["enabled"]) and node["status"] == "active"
|
||||
|
||||
|
||||
def desired_state(node):
|
||||
import db as dbmod
|
||||
|
||||
active = dbmod.list_active_subscriptions(node=node["code"])
|
||||
wanted = {s["uuid"]: s["uuid"] for s in active}
|
||||
nodes_by_code = {n["code"]: n for n in dbmod.list_nodes()}
|
||||
entry_chains = []
|
||||
exit_nodes = {}
|
||||
relay_wanted = {}
|
||||
for chain in dbmod.list_chains(enabled_only=True):
|
||||
entry = nodes_by_code.get(chain["entry_node"])
|
||||
exit_node = nodes_by_code.get(chain["exit_node"])
|
||||
if not entry or not exit_node:
|
||||
continue
|
||||
if not _node_usable(entry) or not _node_usable(exit_node):
|
||||
continue
|
||||
if chain["entry_node"] == node["code"]:
|
||||
entry_chains.append(chain)
|
||||
exit_nodes[chain["exit_node"]] = exit_node
|
||||
if chain["exit_node"] == node["code"] and node["kind"] in ("local", "managed") and chain.get("relay_uuid"):
|
||||
relay_wanted[chain["relay_uuid"]] = chains.relay_email(chain["code"])
|
||||
return wanted, relay_wanted, entry_chains, exit_nodes
|
||||
|
||||
|
||||
def _read_config_text():
|
||||
with open(XRAY_CONFIG_PATH, "r", encoding="utf-8") as f:
|
||||
return f.read()
|
||||
|
||||
|
||||
def _restore_config_text(text):
|
||||
tmp = XRAY_CONFIG_PATH + ".restore.tmp"
|
||||
with open(tmp, "w", encoding="utf-8") as f:
|
||||
f.write(text)
|
||||
os.replace(tmp, XRAY_CONFIG_PATH)
|
||||
subprocess.run(["systemctl", "restart", "xray"], timeout=20)
|
||||
|
||||
|
||||
def _reload_and_verify():
|
||||
subprocess.run(["systemctl", "restart", "xray"], check=True, timeout=20)
|
||||
time.sleep(1)
|
||||
state = subprocess.run(["systemctl", "is-active", "xray"], capture_output=True, text=True).stdout.strip()
|
||||
if state != "active":
|
||||
raise ConfigValidationError("xray не поднялся после применения конфига, вернули старый")
|
||||
|
||||
|
||||
def _port_busy(port):
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
try:
|
||||
sock.bind(("0.0.0.0", int(port)))
|
||||
return False
|
||||
except OSError:
|
||||
return True
|
||||
finally:
|
||||
sock.close()
|
||||
|
||||
|
||||
def _open_firewall(port):
|
||||
subprocess.run(
|
||||
["sh", "-c", f"command -v ufw >/dev/null 2>&1 && ufw allow {int(port)}/tcp || true"],
|
||||
timeout=20,
|
||||
)
|
||||
|
||||
|
||||
def _reconcile_local(wanted, relay_wanted, entry_chains, exit_nodes, apply_chains=True):
|
||||
with _locked():
|
||||
before_text = _read_config_text()
|
||||
cfg = json.loads(before_text)
|
||||
usable = entry_chains
|
||||
skipped = []
|
||||
new_ports = []
|
||||
if apply_chains:
|
||||
busy = {c["port"] for c in entry_chains if _port_busy(c["port"])}
|
||||
usable, skipped = chains.split_busy_chains(cfg, entry_chains, busy)
|
||||
new_ports = chains.new_ports_needed(cfg, usable)
|
||||
changed, problems = chains.sync_config(cfg, wanted, relay_wanted, usable, exit_nodes, apply_chains=apply_chains)
|
||||
if changed:
|
||||
_save(cfg)
|
||||
try:
|
||||
_reload_and_verify()
|
||||
except Exception:
|
||||
_restore_config_text(before_text)
|
||||
raise
|
||||
for port in new_ports:
|
||||
_open_firewall(port)
|
||||
return {"changed": changed, "new_ports": new_ports, "problems": skipped + problems}
|
||||
|
||||
|
||||
def sync_node(node):
|
||||
wanted, relay_wanted, entry_chains, exit_nodes = desired_state(node)
|
||||
if node["kind"] == "managed":
|
||||
import nodeprov
|
||||
reconcile = nodeprov.remote_reconcile
|
||||
args = (node, wanted, relay_wanted, entry_chains, exit_nodes)
|
||||
elif node["kind"] == "local":
|
||||
reconcile = _reconcile_local
|
||||
args = (wanted, relay_wanted, entry_chains, exit_nodes)
|
||||
else:
|
||||
return {"changed": False, "new_ports": [], "problems": []}
|
||||
try:
|
||||
return reconcile(*args)
|
||||
except Exception as first_error:
|
||||
if not entry_chains:
|
||||
raise
|
||||
result = reconcile(*args, apply_chains=False)
|
||||
result["problems"].append(f"цепочки не применились, клиенты синхронизированы: {first_error}")
|
||||
return result
|
||||
|
||||
|
||||
def sync_from_db():
|
||||
import db as dbmod
|
||||
|
||||
expired = dbmod.deactivate_expired()
|
||||
active = dbmod.list_active_subscriptions(node="de1")
|
||||
active_by_id = {s["uuid"]: s for s in active}
|
||||
node = dbmod.get_node("de1")
|
||||
result = sync_node(node)
|
||||
wanted = desired_state(node)[0]
|
||||
return {
|
||||
"removed_expired": len(expired), "active_now": len(wanted),
|
||||
"reloaded": result["changed"], "problems": result["problems"],
|
||||
}
|
||||
|
||||
with _locked():
|
||||
cfg = _load()
|
||||
changed = False
|
||||
for ib in _local_inbounds(cfg):
|
||||
clients = ib["settings"]["clients"]
|
||||
current_ids = {c["id"] for c in clients}
|
||||
if current_ids == set(active_by_id.keys()):
|
||||
continue
|
||||
new_clients = [c for c in clients if c["id"] in active_by_id]
|
||||
existing_ids = {c["id"] for c in new_clients}
|
||||
flow = _TAG_FLOW.get(ib["tag"])
|
||||
for cid, sub in active_by_id.items():
|
||||
if cid not in existing_ids:
|
||||
entry = {"id": cid, "email": cid}
|
||||
if flow:
|
||||
entry["flow"] = flow
|
||||
new_clients.append(entry)
|
||||
ib["settings"]["clients"] = new_clients
|
||||
changed = True
|
||||
|
||||
if changed:
|
||||
_save(cfg)
|
||||
_reload_xray()
|
||||
|
||||
return {"removed_expired": len(expired), "active_now": len(active_by_id), "reloaded": changed}
|
||||
def probe_from_node(node: dict, host: str, port: int):
|
||||
if node["kind"] == "local":
|
||||
return chains.tcp_connect_ms(host, port)
|
||||
if node["kind"] == "managed":
|
||||
import nodeprov
|
||||
return nodeprov.remote_probe(node, host, port)
|
||||
raise ValueError("нода не под управлением панели, замерить с неё нельзя")
|
||||
|
||||
|
||||
def add_client_to_node(node: dict, client_uuid: str, email: str):
|
||||
|
|
@ -266,7 +367,6 @@ def local_node_status() -> dict:
|
|||
|
||||
def sync_all():
|
||||
import db as dbmod
|
||||
import nodeprov
|
||||
|
||||
expired = dbmod.deactivate_expired()
|
||||
results = {}
|
||||
|
|
@ -276,8 +376,15 @@ def sync_all():
|
|||
elif node["kind"] == "managed":
|
||||
active = dbmod.list_active_subscriptions(node=node["code"])
|
||||
try:
|
||||
nodeprov.remote_sync(node, active)
|
||||
results[node["code"]] = {"active_now": len(active), "ok": True}
|
||||
res = sync_node(node)
|
||||
results[node["code"]] = {
|
||||
"active_now": len(active), "ok": True,
|
||||
"changed": res["changed"], "problems": res["problems"],
|
||||
}
|
||||
except Exception as e:
|
||||
results[node["code"]] = {"active_now": len(active), "ok": False, "error": str(e)}
|
||||
return {"removed_expired": len(expired), "nodes": results}
|
||||
reloaded = any(r.get("changed") or r.get("reloaded") for r in results.values())
|
||||
return {
|
||||
"removed_expired": len(expired), "active_now": len(dbmod.list_active_subscriptions()),
|
||||
"reloaded": reloaded, "nodes": results,
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue