services v1.0.0 [extra]
Управление службами: Windows/Linux/Docker/локальные юниты, каталог, действия.
| version | date | commit | файлов |
|---|---|---|---|
| 1.0.0 | 2026-10-07 | 8c15e4d68659 | 23 |
README
# services
Управление службами на всех хостах из одного места: что есть, что запущено, старт/стоп/рестарт.
Тип: Модуль. Категория: `extra`. Зависимостей нет.
## Что это
Службы живут на разных машинах и разных init-системах — Windows (`sc`), Linux (`systemctl` по SSH), Docker-контейнеры, локальные юниты. Спрашивать хочется одинаково, поэтому ядро держит реестр, а провайдеры отвечают каждый за свой тип:
- `services_windows` — службы Windows (`sc query/start/stop`, автозапуск через `sc config`);
- `services_linux` — юниты удалённых Linux-машин (`systemctl` по SSH, действия через sudo);
- `services_docker` — контейнеры (`docker ps/start/stop/restart`, включая healthy);
- `services_local` — юниты своей же машины (`systemctl` напрямую, без SSH).
Провайдеры видят только свой вайтлист: остальные службы хоста боту недоступны. Каждое действие пишется в аудит, статусы собираются best-effort (упавший хост не роняет весь список).
В кабинете — сводная таблица служб всех хостов со статусами, в Telegram — кнопка сервисов в меню (группировка по хостам).
## Каталог
Право `service_catalog` (старт/стоп/рестарт не даёт): все доступные службы хостов — из кеша `data/services/discovery_cache.json`, обновление только вручную по кнопке (по хосту или всё сразу). Чекбокс «вкл» — служба в общем списке, выкл — скрыта.
Списки хранятся в `data/services/hosts.json` (вне git, один файл на все типы: `type/alias/endpoint/services/...`), сохранение из кабинета — без рестарта (`hosts_set` + `reload` провайдеров). Нет `hosts.json` — провайдеры работают как раньше: один инстанс из своего `.env`.
Подключения — одна таблица перед каталогом: хост, тип, селект точки → привязка, кеш, действия (🔄 опрос, 💾 сохранить, 🗑 удалить точку). Новые точки создаются там же: windows/linux/docker → ssh-точка, local — новый хост без точки. SSH-пароль/ключ живёт в точке (`data/env/outbound.env`), здесь — только вайтлисты и sudo-пароли для удалённых команд. С ключом (`address.key`) пароль можно не вводить.
## Устройство (для разработчиков)
Провайдеры встроены в ядро: движки — `services/backends/` (`docker.py/linux.py/local.py/windows.py` + базовый `ProviderModule`), конфиги — `backends/<type>.schema.json`, примеры — `backends/<type>.env.example`. Значения читаются из `data/env/services_<type>.env` (имена сохранены). API-имена провайдеров (`services_docker.*` и т.д.) сохранены ради совместимости кабинета и внешних вызовов. Контракт провайдера — пять API-методов (`services_info/service_status/service_action/discover/reload`), ядро вызывает их по имени.
Отдельных модулей `services_*` больше нет (релиз 1.0.0): при первом старте ядро само сносит их каталоги с бекапом в `data/backups/modules-legacy/` и помечает записи в modstate как missing — обновление приезжает обычной заменой файлов.
## Возможности
- `register {provider: str}` → регистрация провайдера (зовут services_* в своём setup)
- `hosts_for {svc_type}` → инстансы типа из hosts.json (валидированные; [] — fallback на .env)
- `hosts_all {}` → весь hosts.json — для кабинета
- `hosts_set {entries, actor = ""}` → перезаписать hosts.json (атомарно) + reload провайдеров
- `discover {host = "", actor = ""}` → опросить хосты, перезаписать кеш каталога
- `catalog {query, type_, host, checked, page, per_page}` → каталог из кеша: фильтры + пагинация
- `endpoint_add {name, svc_type, address, secret, via?, note?}` → создать точку outbound (ssh) с пробой подключения
- `endpoint_delete {name}` → удалить точку; ссылки инстансов сбрасываются (reload)
- `host_add {alias, svc_type, endpoint?, sudo_pass?}` → новый хост в hosts.json (local — без точки, services пустой)
- `providers {}` → провайдеры и их инстансы
- `list {}` → плоский список служб (имена + host/type)
- `status {}` → статусы всех служб (best-effort)
- `action {host: str, name: str, act: start|stop|restart}` → действие с маршрутизацией по host
## Права
- `service_manage` — управление службами (таблица статусов, старт/стоп/рестарт)
- `service_catalog` — каталог (hosts/точки/кеш), старт/стоп/рестарт не даёт
## Telegram
- Кнопка сервисов в меню (колбэки `service_*`, группировка по хостам); команд нет
## Web
- `GET /api/admin/services` — статусы служб всех хостов
- `POST /api/admin/service` — действие (`start|stop|restart`)
- `GET /api/admin/services_catalog` — каталог из кеша
- `POST /api/admin/services_discover` — обновить кеш (в фоне)
- `GET/POST /api/admin/services_hosts` — инстансы и чекбоксы
- `POST /api/admin/services_endpoint` — создать точку outbound
- `POST /api/admin/services_endpoint_delete` — удалить точку
- `POST /api/admin/services_host_add` — добавить хост
## Конфигурация (`config.schema.json`)
- обязательных и опциональных нет (настройки провайдеров — в их `.env`, см. ниже)
Манифест
{
"name": "services",
"version": "1.0.0",
"category": "extra",
"description": "Управление службами: Windows/Linux/Docker/локальные юниты, каталог, действия.",
"requires": [],
"rights": [],
"api_methods": [],
"web_routes": [],
"commit": "8c15e4d68659850405ecb08860eebe6d3a4d2f45",
"updated": "2026-10-07 10:00:57 +0300",
"integrity": "sha256:5cc3fb456522190407a236dfad7ecc6988294008de299ed8df913166e27c28a2",
"files": "[23 файлов — см. вкладки ниже]",
"note": "",
"wdesc": "",
"wbody": "",
"_extra": {
"api": "",
"events_emitted": "[]",
"events_subscribed": "[]",
"audit_allow": "",
"shared_routes": "",
"config_schema": "config.schema.json",
"env_example": ".env.example"
}
} Файлы и исходники
Дерево файлов
- · корень
- 6.8 КБ
- 0.0 КБ
- 6.0 КБ
- 34.2 КБ
- 6.2 КБ
- 0.8 КБ
- 11.4 КБ
- backends/
- 0.6 КБ
- 0.7 КБ
- 3.0 КБ
- 0.1 КБ
- 0.7 КБ
- 4.8 КБ
- 0.1 КБ
- 0.4 КБ
- 3.7 КБ
- 0.1 КБ
- 11.6 КБ
- 0.7 КБ
- 7.2 КБ
- 0.2 КБ
- web/
- 5.0 КБ
- 14.2 КБ
Предпросмотр
Выберите файл в дереве выше — код откроется здесь.
Все исходники (.py) одним списком
backends/__init__.py 0.6 КБ
"""services.backends — движки провайдеров служб (без loader-обвязки).
docker/linux/windows — команды по SSH через runner, инжектируемый
конструктором (cmd, timeout) -> (rc, output); local — subprocess напрямую.
Вайщность вайтлистов и экранирование имён — внутри классов (см. _check).
Чистые парсеры (parse_sc_names, parse_systemctl_names) — для каталога
и тестов, без runner/subprocess.
"""
backends/docker.py 3.0 КБ
"""services.backends.docker — Docker-хост: ps/start/stop/restart по SSH.
Транспорт — только через outbound (точка kind=ssh, ssh.exec): runner
инжектится конструктором (cmd, timeout) -> (rc, output). Локального SSH
здесь нет. Контейнеры — только вайтлист DOCKER_CONTAINERS.
"""
import logging
import re as _re
logger = logging.getLogger(__name__)
def split_list(value):
return [x.strip() for x in (value or "").replace(";", ",").split(",") if x.strip()]
class DockerServices:
def __init__(self, runner, containers):
self._run = runner
self.containers = list(containers or [])
def _check(self, name):
if name not in self.containers:
raise ValueError("контейнер %s не в вайтлисте" % name)
# имя уходит в удалённый shell (docker ...): только символы
# имён контейнеров, метасимволы запрещены даже для вайтлиста.
if not _re.fullmatch(r"[A-Za-z0-9_.-]{1,128}", name or ""):
raise ValueError("недопустимые символы в имени контейнера")
def statuses(self, max_chars=0):
"""{name: (up, status)} по docker ps (все, дальше фильтр по вайтлисту)."""
params_max = max_chars or 200000
rc, out = self._run(
"docker ps -a --format '{{.Names}}|{{.Status}}'", 30,
max_chars=params_max)
if rc != 0:
raise RuntimeError(out.strip()[:300] or ("exit %d" % rc))
res = {}
for line in out.splitlines():
if "|" not in line:
continue
nm, st = line.split("|", 1)
nm, st = nm.strip(), st.strip()
if nm:
res[nm] = (st.upper().startswith("UP"), st)
return res
def discover(self):
"""Все контейнеры хоста (без фильтра вайтлиста) — для каталога."""
return sorted(self.statuses())
def status(self, name):
self._check(name)
try:
table = self.statuses()
except Exception as e: # noqa: BLE001
return {"name": name, "up": False, "detail": str(e)[:200]}
if name not in table:
return {"name": name, "up": False, "detail": "контейнер не найден"}
up, st = table[name]
return {"name": name, "up": up, "detail": st[:200]}
def action(self, name, act):
self._check(name)
if act not in ("start", "stop", "restart"):
return {"ok": False, "error": "недопустимое действие"}
rc, out = self._run("docker %s %s" % (act, name), 90)
out = out.strip()[:500]
if rc != 0:
return {"ok": False, "error": out or ("exit %d" % rc)}
return {"ok": True, "output": out or "OK"}
backends/linux.py 4.8 КБ
"""services.backends.linux — удалённый Linux-хост: systemctl по SSH.
Транспорт — только через outbound (точка kind=ssh, ssh.exec): runner
инжектится конструктором (cmd, timeout) -> (rc, output). Локального SSH
здесь нет. Статус — без sudo. Действия — `echo PASS | sudo -S systemctl ...`
(sudo-пароль — из конфига провайдера, используется внутри удалённой
команды). Юниты — только вайтлист LINUX_SERVICES.
"""
import logging
import re as _re
logger = logging.getLogger(__name__)
def split_list(value):
return [x.strip() for x in (value or "").replace(";", ",").split(",") if x.strip()]
def _q(s):
"""Экранирование для одинарных кавычек в shell."""
return str(s or "").replace("'", "'\\''")
def parse_systemctl_names(out):
"""Имена юнитов из list-units/list-unit-files (первый токен, без .service).
Чистая функция (без runner) — для каталога и тестов.
"""
names = []
for line in (out or "").splitlines():
line = line.strip()
if not line:
continue
first = line.split()[0]
if first.endswith(".service"):
first = first[:-len(".service")]
names.append(first)
seen = set()
res = []
for n in names:
if n and n not in seen:
seen.add(n)
res.append(n)
return res
class LinuxServices:
def __init__(self, runner, services, sudo_password=""):
self._run = runner
self.services = list(services or [])
self.sudo_password = sudo_password or ""
def _check(self, name):
if name not in self.services:
raise ValueError("юнит %s не в вайтлисте" % name)
# имя уходит в удалённый shell (systemctl ...): только символы
# имён юнитов systemd, метасимволы запрещены даже для вайтлиста.
if not _re.fullmatch(r"[A-Za-z0-9_.:@-]{1,128}", name or ""):
raise ValueError("недопустимые символы в имени юнита")
def discover(self):
"""Все юниты хоста (без фильтра вайтлиста, read-only без sudo)."""
names = []
for cmd in ("systemctl list-units --type=service --all "
"--no-legend --no-pager",
"systemctl list-unit-files --type=service "
"--no-legend --no-pager"):
rc, out = self._run(cmd, 90, max_chars=1000000)
if rc != 0:
raise RuntimeError(out.strip()[:300] or ("exit %d" % rc))
names.extend(parse_systemctl_names(out))
return sorted(set(names))
def status(self, name):
self._check(name)
rc, out = self._run(
"systemctl --no-pager is-active %s" % name, 30)
state = out.strip().splitlines()[0].strip() if out.strip() else ""
if rc != 0 and not state:
return {"name": name, "up": False,
"detail": out.strip()[:200] or ("exit %d" % rc)}
return {"name": name, "up": state == "active", "detail": state[:200]}
def action(self, name, act):
self._check(name)
if act not in ("start", "stop", "restart"):
return {"ok": False, "error": "недопустимое действие"}
if not self.sudo_password:
return {"ok": False, "error": "не задан LINUX_SUDO_PASS"}
remote = "echo '%s' | sudo -S -p '' systemctl %s %s" % (
_q(self.sudo_password), act, name)
rc, out = self._run(remote, 90)
out = out.strip()[:500]
if rc != 0:
return {"ok": False, "error": out or ("exit %d" % rc)}
import time as _time
if act == "stop":
_time.sleep(3)
final = self._wait_state(name, act != "stop")
up = bool(final.get("up"))
if act != "stop" and not up:
return {"ok": False,
"error": "не поднялся: %s" % (final.get("detail") or "?")}
return {"ok": True, "output": final.get("detail") or "OK", "up": up}
def _wait_state(self, name, want_up, timeout=120):
import time as _time
deadline = _time.monotonic() + max(5, timeout)
last = None
while _time.monotonic() < deadline:
try:
st = self.status(name)
except Exception: # noqa: BLE001
st = {"up": False, "detail": ""}
last = st
if bool(st.get("up")) == bool(want_up):
return st
_time.sleep(5)
return last or {"name": name, "up": False, "detail": "таймаут ожидания"}
backends/local.py 3.7 КБ
"""services.backends.local — локальные systemd-юниты.
Статус — обычный systemctl (без sudo). Действия — sudo -S, пароль из .env
модуля (LOCAL_SUDO_PASS). Юниты — только вайтлист LOCAL_SERVICES.
"""
import logging
import subprocess
logger = logging.getLogger(__name__)
def split_list(value):
return [x.strip() for x in (value or "").replace(";", ",").split(",") if x.strip()]
def parse_systemctl_names(out):
"""Имена юнитов из list-units/list-unit-files (первый токен, без .service).
Чистая функция (без subprocess) — для каталога и тестов.
"""
names = []
for line in (out or "").splitlines():
line = line.strip()
if not line:
continue
first = line.split()[0]
if first.endswith(".service"):
first = first[:-len(".service")]
names.append(first)
seen = set()
res = []
for n in names:
if n and n not in seen:
seen.add(n)
res.append(n)
return res
class LocalServices:
def __init__(self, services, sudo_password=""):
self.services = list(services or [])
self.sudo_password = sudo_password or ""
def _check(self, name):
if name not in self.services:
raise ValueError("юнит %s не в вайтлисте" % name)
def discover(self):
"""Все юниты локального хоста (без фильтра вайтлиста, без sudo)."""
names = []
for cmd in (["systemctl", "list-units", "--type=service", "--all",
"--no-legend", "--no-pager"],
["systemctl", "list-unit-files", "--type=service",
"--no-legend", "--no-pager"]):
try:
r = subprocess.run(cmd, capture_output=True, text=True,
timeout=60)
except Exception as e: # noqa: BLE001
raise RuntimeError(str(e)[:300])
if r.returncode != 0:
raise RuntimeError((r.stderr or "").strip()[:300]
or ("exit %d" % r.returncode))
names.extend(parse_systemctl_names(r.stdout or ""))
return sorted(set(names))
def status(self, name):
self._check(name)
try:
r = subprocess.run(["systemctl", "is-active", name],
capture_output=True, text=True, timeout=15)
except Exception as e: # noqa: BLE001
return {"name": name, "up": False, "detail": str(e)[:200]}
state = (r.stdout or "").strip()
return {"name": name, "up": state == "active", "detail": state[:200]}
def action(self, name, act):
self._check(name)
if act not in ("start", "stop", "restart"):
return {"ok": False, "error": "недопустимое действие"}
if not self.sudo_password:
return {"ok": False,
"error": "не задан LOCAL_SUDO_PASS"}
try:
r = subprocess.run(
["sudo", "-S", "-p", "", "systemctl", act, name],
input=self.sudo_password + "\n", capture_output=True,
text=True, timeout=90)
except subprocess.TimeoutExpired:
return {"ok": False, "error": "таймаут systemctl"}
err = (r.stderr or "").strip()
if r.returncode != 0:
return {"ok": False, "error": err[:500] or ("exit %d" % r.returncode)}
st = self.status(name)
return {"ok": True,
"output": ("active" if st["up"] else st["detail"]) or "OK",
"up": st["up"]}
backends/provider.py 11.6 КБ
"""services.backends.provider — базовый класс провайдера служб.
Выносит общий код четырёх встроенных провайдеров (docker/linux/windows/
local): регистрация API (services_info/service_status/service_action/
discover/reload), инстансы из hosts.json/.env, резолв имён, discover
best-effort. Отличия задаются атрибутами класса (см. _EmbeddedProvider
в modules/services/module.py):
MOD — API-имя ("services_docker"); TYPE — тип ("docker").
ENGINE — класс движка (DockerServices/...); ENGINE_ARGS — как строить
движок из конфига инстанса: "containers" (docker),
"services+autostart" (windows), "services+sudo" (linux),
"services+sudo-local" (local, без outbound).
NEED_OUTBOUND — True (ssh через outbound) / False (local).
CONSUMER_DESC — описание потребителя для outbound.consumer_ensure.
LEGACY_* — ключи .env для одиночного инстанса без hosts.json:
LEGACY_ALIAS_KEY, LEGACY_SERVICES_KEY, LEGACY_EXTRA — dict
{attr: env_key} для autostart/sudo ("" — пропустить).
UNIT_NOUN — слово для ошибок резолва ("контейнер"/"юнит"/"сервис").
HEALTH_COUNT — ключ health-счётчика ("containers"/"units"/"services").
"""
import os
from core.base_module import BaseModule
class ProviderModule(BaseModule):
MOD = ""
TYPE = ""
ENGINE = None
ENGINE_ARGS = "containers"
NEED_OUTBOUND = True
CONSUMER_DESC = ""
LEGACY_ALIAS_KEY = ""
LEGACY_ALIAS_DEFAULT = ""
LEGACY_SERVICES_KEY = ""
LEGACY_EXTRA = {}
UNIT_NOUN = "служба"
HEALTH_COUNT = "services"
def setup(self, ctx) -> None:
mod_dir = os.path.dirname(os.path.abspath(__file__))
self.config = self._provider_config(ctx, mod_dir)
self.ctx = ctx
self._load()
ctx.api.register(self.MOD, "services_info", self.info)
ctx.api.register(self.MOD, "service_status", self.status)
ctx.api.register(self.MOD, "service_action", self.action)
ctx.api.register(self.MOD, "discover", self.discover)
ctx.api.register(self.MOD, "reload", self.reload)
def _provider_config(self, ctx, mod_dir):
"""Конфиг провайдера: data/env/<MOD>.env поверх схемы бэкенда.
Схема — backends/<type>.schema.json (дефолты), значения —
data/env/services_<type>.env (имена сохранены: существующие
установки ничего не теряют). mod_dir как module_dir не годится
(там нет config.schema.json провайдера) — собираем ModuleConfig
напрямую.
"""
import json as _json
from core.config import ModuleConfig
schema = {"required": [], "optional": {}}
sp = os.path.join(os.path.dirname(os.path.abspath(__file__)),
"%s.schema.json" % self.TYPE)
try:
with open(sp, encoding="utf-8") as f:
schema = _json.load(f)
except Exception: # noqa: BLE001
pass
overrides = {}
try:
if getattr(ctx, "runtime", None) is not None:
overrides = ctx.runtime.get_overrides(self.MOD) or {}
except Exception: # noqa: BLE001
pass
return ModuleConfig(self.MOD, mod_dir, schema,
overrides=overrides,
env_dir=os.path.join(ctx.data_dir, "env"))
# --- инстансы ---
def _split(self, value):
return [x.strip() for x in (value or "").replace(";", ",").split(",")
if x.strip()]
def _legacy_cfg(self):
"""Один инстанс из .env (нет hosts.json): (alias, consumer, kwargs)."""
alias = (self.config.get(self.LEGACY_ALIAS_KEY, "") or "").strip() \
or self.LEGACY_ALIAS_DEFAULT
kw = {"services": self._split(
self.config.get(self.LEGACY_SERVICES_KEY, ""))}
for attr, env_key in (self.LEGACY_EXTRA or {}).items():
if attr == "autostart":
kw[attr] = self._split(self.config.get(env_key, ""))
else:
kw[attr] = (self.config.get(env_key, "") or "")
if self.NEED_OUTBOUND:
return [(alias, self.MOD, kw)]
return [(alias, None, kw)]
def _make_prov(self, runner, kw):
args = self.ENGINE_ARGS
if args == "containers":
return self.ENGINE(runner, kw["services"])
if args == "services+autostart":
return self.ENGINE(runner, kw["services"],
kw.get("autostart") or [])
if args == "services+sudo":
return self.ENGINE(runner, kw["services"],
kw.get("sudo_pass") or "")
if args == "services+sudo-local":
return self.ENGINE(kw["services"],
kw.get("sudo_pass") or "")
raise ValueError("неизвестный ENGINE_ARGS: %s" % args)
def _load(self):
"""(Пере)собрать инстансы из hosts.json/.env (setup + reload)."""
ctx = self.ctx
try:
cfgs = ctx.api.call("services.hosts_for", self.TYPE)
except Exception as e: # noqa: BLE001
from core.errors import ConfigError
raise ConfigError("%s: %s" % (self.MOD, e))
inst = []
if not cfgs:
for alias, consumer, kw in self._legacy_cfg():
if self.NEED_OUTBOUND:
self._ensure_consumer(consumer)
prov = self._make_prov(self._make_runner(consumer), kw)
else:
prov = self._make_prov(None, kw)
inst.append({"alias": alias,
**({"consumer": consumer}
if consumer else {}),
"prov": prov})
else:
for ic in cfgs:
consumer = "%s:%s" % (self.MOD, ic["alias"])
if self.NEED_OUTBOUND:
self._ensure_consumer(consumer)
if ic.get("endpoint"):
try:
ctx.api.call("outbound.bind", consumer,
ic["endpoint"],
"services:autobind")
except Exception as e: # noqa: BLE001
ctx.get_logger(self.MOD).warning(
"авто-привязка %s -> %s не удалась: %s "
"(создайте точку в кабинете outbound)",
consumer, ic["endpoint"], e)
kw = {"services": ic["services"]}
if self.ENGINE_ARGS == "services+autostart":
kw["autostart"] = ic.get("autostart") or []
elif self.ENGINE_ARGS == "services+sudo":
kw["sudo_pass"] = ic.get("sudo_pass") or ""
prov = self._make_prov(self._make_runner(consumer),
kw)
else:
kw = {"services": ic["services"]}
if self.ENGINE_ARGS == "services+sudo-local":
kw["sudo_pass"] = ic.get("sudo_pass") or ""
prov = self._make_prov(None, kw)
inst.append({"alias": ic["alias"],
**({"consumer": consumer}
if self.NEED_OUTBOUND else {}),
"prov": prov})
self._inst = inst
def _ensure_consumer(self, consumer):
try:
self.ctx.api.call("outbound.consumer_ensure", consumer,
"ssh.exec", self.CONSUMER_DESC)
except Exception as e: # noqa: BLE001
from core.errors import ConfigError
raise ConfigError("%s: нет шлюза outbound: %s" % (self.MOD, e))
def reload(self):
"""Перечитать hosts.json без рестарта (зовёт services.hosts_set)."""
self._load()
return {"ok": True, "hosts": [i["alias"] for i in self._inst]}
def _make_runner(self, consumer):
def _run(cmd, timeout=30, max_chars=0):
"""(rc, output) через outbound ssh.exec; без привязки — исключение."""
params = {"cmd": cmd, "timeout": timeout}
if max_chars:
params["max_chars"] = max_chars
res = self.ctx.api.call("outbound.call", consumer, "ssh.exec",
params) or {}
if "rc" not in res:
raise RuntimeError(res.get("error") or "outbound: нет ответа")
return int(res.get("rc") or 0), res.get("output") or ""
return _run
# --- контракт ядра ---
def _names(self, prov):
cont = getattr(prov, "containers", None)
if cont is not None:
return list(cont)
return list(prov.services)
def _resolve(self, name, host=""):
if host:
for i in self._inst:
if i["alias"] == (host or ""):
return i
raise ValueError("хост %s не найден" % host)
hits = [i for i in self._inst if name in self._names(i["prov"])]
if len(hits) == 1:
return hits[0]
if not hits:
raise ValueError("%s %s не в вайтлисте" % (self.UNIT_NOUN,
name))
raise ValueError("%s %s есть на нескольких хостах "
"(укажите host)" % (self.UNIT_NOUN, name))
def info(self):
"""Список инстансов [{host, type, services}] (вайтлисты)."""
return [{"host": i["alias"], "type": self.TYPE,
"services": self._names(i["prov"])}
for i in self._inst]
def status(self, name, host=""):
"""Статус службы; host — для инстанса."""
return self._resolve(name, host)["prov"].status(name)
def action(self, name, act, host=""):
"""Действие; host — для инстанса."""
res = self._resolve(name, host)["prov"].action(name, act)
if not res.get("ok"):
raise RuntimeError(res.get("error") or "ошибка")
return res
def discover(self, host=""):
"""Все службы хостов (без фильтра, best-effort) — для каталога."""
out = []
for i in self._inst:
if host and i["alias"] != host:
continue
try:
svc = i["prov"].discover()
out.append({"host": i["alias"], "type": self.TYPE,
"services": svc})
except Exception as e: # noqa: BLE001
out.append({"host": i["alias"], "type": self.TYPE,
"services": [],
"error": "%s: %s" % (type(e).__name__, e)})
if host and not out:
raise ValueError("хост %s не найден" % host)
return out
def health(self) -> dict:
return {"ok": True, "module": self.name,
"hosts": [i["alias"] for i in self._inst],
self.HEALTH_COUNT: sum(len(self._names(i["prov"]))
for i in self._inst)}
backends/windows.py 7.2 КБ
"""services.backends.windows — Windows-хост: sc query/start/stop/config по SSH.
Транспорт — только через outbound (точка kind=ssh, ssh.exec): runner
инжектится конструктором (cmd, timeout) -> (rc, output). Локального SSH
здесь нет. Сервисы — только вайтлист WIN_SERVICES (все подряд не трогаем).
"""
import logging
import re as _re
logger = logging.getLogger(__name__)
def split_list(value):
return [x.strip() for x in (value or "").replace(";", ",").split(",") if x.strip()]
def parse_sc_names(out):
"""Имена служб из `sc query type= service state= all`.
Чистая функция (без runner) — для каталога и тестов.
Запись sc отделяется пустой строкой, первое поле — всегда
SERVICE_NAME. Метка поля может быть побита кодировкой (cp866
через SSH: `���_�㦡�: llama-qwen36-moe`), поэтому матчим сначала
строго, затем — первую строку записи вида `xxx: <токен>`.
"""
names = []
for block in _re.split(r"\r?\n\s*\r?\n", out or ""):
lines = [ln for ln in block.splitlines() if ln.strip()]
if not lines:
continue
m = _re.match(r"\s*SERVICE_NAME\s*:\s*(\S+)", lines[0], _re.I)
if m:
names.append(m.group(1).strip())
continue
# побитая метка: SERVICE_NAME всегда с нулевого отступа,
# поля записи — с отступом; требуем первую строку без отступа
# вида `xxx: <токен>` + вторую строку с двоеточием
if lines[0][:1] in (" ", "\t"):
continue
m = _re.match(r"\S[^:]*:\s*(\S+)\s*$", lines[0])
if m and len(lines) > 1 and ":" in lines[1]:
names.append(m.group(1).strip())
seen = set()
res = []
for n in names:
if n not in seen:
seen.add(n)
res.append(n)
return res
class WinServices:
def __init__(self, runner, services, autostart):
self._run = runner
self.services = list(services or [])
self.autostart = set(autostart or [])
def _check(self, name):
if name not in self.services:
raise ValueError("сервис %s не в вайтлисте" % name)
# имя уходит в удалённый shell (sc query/start ...): метасимволы
# запрещены даже для вайтлиста (пробелы в именах допустимы).
if _re.search(r"[;|&$`\"'\\<>!*?~#%()\[\]{}^]", name or "") or \
"\n" in (name or "") or "\r" in (name or ""):
raise ValueError("недопустимые символы в имени сервиса")
def discover(self):
"""Все службы хоста (без фильтра вайтлиста) — для каталога."""
rc, out = self._run("sc query type= service state= all", 90,
max_chars=1000000)
if rc != 0:
raise RuntimeError(out.strip()[:300] or ("exit %d" % rc))
return parse_sc_names(out)
def status(self, name):
self._check(name)
rc, out = self._run("sc query %s" % name, 30)
up = "RUNNING" in out.upper()
detail = ""
for line in out.splitlines():
if "STATE" in line.upper():
detail = line.strip()
break
# cp866-мусор в STATE-строке — схлопываем до сути
m = _re.search(r"(RUNNING|STOPPED)", (detail or out).upper())
return {"name": name, "up": up,
"detail": (m.group(1).capitalize() if m else detail)[:200]
or out.strip()[:200]}
def _wait_state(self, name, want_up, timeout=120):
"""Ждём фактического состояния (медленный старт с загрузкой модели)."""
import time as _time
deadline = _time.monotonic() + max(5, timeout)
last = None
while _time.monotonic() < deadline:
try:
st = self.status(name)
except Exception: # noqa: BLE001
st = {"up": False, "detail": ""}
last = st
if bool(st.get("up")) == bool(want_up):
return st
_time.sleep(5)
return last or {"name": name, "up": False, "detail": "таймаут ожидания"}
def action(self, name, act):
self._check(name)
if act == "start":
# сначала автозапуск, потом старт (иначе 1058 на disabled)
cmds = []
if name in self.autostart:
cmds.append("sc config %s start= auto" % name)
cmds.append("sc start %s" % name)
elif act == "stop":
cmds = ["sc stop %s" % name]
if name in self.autostart:
cmds.append("sc config %s start= disabled" % name)
elif act == "restart":
cmds = ["sc stop %s" % name, "__sleep__",
"sc start %s" % name]
else:
return {"ok": False, "error": "недопустимое действие"}
outs = []
for c in cmds:
if c == "__sleep__":
import time as _time
_time.sleep(5) # дать службе корректно завершиться
continue
rc, out = self._run(c, 60)
outs.append(out.strip()[:500])
if rc != 0:
# 1058 (отключена) — повторяем config+start; 1056 (гаснет
# старый инстанс) — ждём 15с и повторяем старт
if act == "start" and c.startswith("sc start"):
import time as _time
if "1058" in (out or ""):
rc2, _o2 = self._run(
"sc config %s start= auto" % name, 60)
if rc2 == 0:
rc, out = self._run(c, 60)
outs.append(out.strip()[:500])
elif "1056" in (out or ""):
_time.sleep(15)
rc, out = self._run(c, 90)
outs.append(out.strip()[:500])
if rc != 0:
return {"ok": False,
"error": " ; ".join(o for o in outs if o) or ("exit %d" % rc)}
if act == "stop":
import time as _time
_time.sleep(3) # отстой SCM перед финальным опросом
final = self._wait_state(name, act != "stop")
up = bool(final.get("up"))
base = " ; ".join(o for o in outs if o) or "OK"
if act != "stop" and not up:
return {"ok": False,
"error": "%s (не поднялась: %s)" % (base, final.get("detail") or "?")}
return {"ok": True, "output": base, "up": up}
module.py 34.2 КБ
"""services.Module — ядро управления службами + встроенные провайдеры.
Встроенные провайдеры (backends/: windows/linux/docker/local) поднимаются
ядром в setup и отдают контракт:
<name>.services_info() -> [{host, type, services}] (список инстансов;
одиночный dict тоже принимается для совместимости)
<name>.service_status(name, host="") -> {"name","up":bool,"detail":""}
<name>.service_action(name, act, host="") -> {"ok","output"/"error"}
<name>.discover(host="") -> [{host, type, services}] (все доступные,
без фильтра вайтлиста — для каталога)
API-имена провайдеров (services_docker.* и т.д.) сохранены ради
совместимости кабинета и внешних вызовов, отдельных модулей больше нет.
services.register оставлен для сторонних провайдеров.
Мульти-инстансы: data/services/hosts.json (вне git), один файл на все
типы: [{"type","alias","endpoint","services",...}]. Нет файла — провайдер
работает как раньше (один инстанс из data/env/services_<type>.env).
Привязка инстанса к точке outbound — по полю endpoint (auto-bind).
Ядро агрегирует: общий список (имена + host/service), статусы, маршрутизация
действий. Пароли — только в .env провайдеров (modsafe 🔑).
"""
import json
import logging
import os
import re as _re
import time as _time
from core.base_module import BaseModule
logger = logging.getLogger(__name__)
#: Типы инстансов в data/services/hosts.json (порядок — для сортировки каталога).
HOSTS_TYPES = ("windows", "linux", "docker", "local")
_ALIAS_RE = _re.compile(r"^[A-Za-z0-9_.-]{1,64}$")
#: Пагинация каталога по умолчанию/максимум.
CATALOG_DEFAULT_PER_PAGE = 50
CATALOG_MAX_PER_PAGE = 500
class Module(BaseModule):
def setup(self, ctx) -> None:
mod_dir = os.path.dirname(os.path.abspath(__file__))
self.config = ctx.module_config("services", mod_dir)
self.ctx = ctx
self._providers = {} # mod_name -> True (порядок регистрации)
self._cleanup_legacy_shims()
ctx.api.register("services", "register", self.api_register)
ctx.api.register("services", "hosts_for", self.api_hosts_for)
ctx.api.register("services", "hosts_all", self.api_hosts_all)
ctx.api.register("services", "hosts_set", self.api_hosts_set)
ctx.api.register("services", "discover", self.api_discover)
ctx.api.register("services", "catalog", self.api_catalog)
ctx.api.register("services", "endpoint_add", self.api_endpoint_add)
ctx.api.register("services", "endpoint_delete",
self.api_endpoint_delete)
ctx.api.register("services", "host_add", self.api_host_add)
ctx.api.register("services", "providers", self.api_providers)
ctx.api.register("services", "list", self.api_list)
ctx.api.register("services", "status", self.api_status)
ctx.api.register("services", "action", self.api_action)
self._embedded = {}
for mod_name, cls in _EMBEDDED_CLASSES.items():
prov = cls(mod_name, mod_dir, ctx)
self._embedded[mod_name] = prov
self._providers.setdefault(mod_name, True)
logger.info("services: встроенный провайдер %s", mod_name)
def health(self) -> dict:
out = {"ok": True, "module": self.name,
"providers": sorted(self._providers)}
try:
hosts = {}
for name, prov in (getattr(self, "_embedded", None) or {}).items():
hosts[name] = [i["alias"] for i in prov._inst]
out["hosts"] = hosts
except Exception: # noqa: BLE001
pass
return out
def _cleanup_legacy_shims(self):
"""Удалить каталоги services_* (релиз 1.0.0: провайдеры встроены).
Обновление через маркет/кабинет привозит только новые файлы —
старые каталоги services_{docker,linux,local,windows} остались бы
лежать и loader поднял бы их шимы дублями API (register падает
ValueError). Каталог-шим детерминированно мусор: удаляем с бекапом
в data/backups/modules-legacy/ (на случай ручных правок внутри).
"""
import shutil as _shutil
import time as _time
mods_dir = os.path.dirname(os.path.dirname(
os.path.abspath(__file__)))
for name in ("services_docker", "services_linux",
"services_local", "services_windows"):
d = os.path.join(mods_dir, name)
if not os.path.isdir(d):
continue
backup = os.path.join(
self.ctx.data_dir, "backups", "modules-legacy",
"%s-%s" % (name, _time.strftime("%Y%m%d_%H%M%S")))
try:
os.makedirs(os.path.dirname(backup), exist_ok=True)
_shutil.copytree(d, backup, symlinks=False,
ignore=_shutil.ignore_patterns(
"__pycache__"))
except Exception as e: # noqa: BLE001
logger.error("services: бекап legacy-шима %s не удался: %s",
name, e)
continue
try:
_shutil.rmtree(d)
logger.warning("services: удалён legacy-шим %s (бекап %s)",
name, backup)
except Exception as e: # noqa: BLE001
logger.error("services: не удался снос legacy-шима %s: %s",
name, e)
# Записи шимов в modstate — в missing (кабинет их больше не покажет
# как модули; их API теперь отдаёт ядро).
try:
from core import modstate as _ms
state = _ms.load(self.ctx.data_dir)
mods = state.get("modules") or {}
touched = False
for name in ("services_docker", "services_linux",
"services_local", "services_windows"):
rec = mods.get(name)
if isinstance(rec, dict) and not rec.get("missing"):
rec["missing"] = True
touched = True
if touched:
_ms.save(self.ctx.data_dir, state)
except Exception as e: # noqa: BLE001
logger.warning("services: modstate legacy-шимы не помечены: %s",
e)
# --- инстансы (data/services/hosts.json) ---
def _hosts_path(self):
return os.path.join(self.ctx.data_dir, "services", "hosts.json")
def _cache_path(self):
return os.path.join(self.ctx.data_dir, "services",
"discovery_cache.json")
@staticmethod
def _validate_hosts(data):
"""Валидированный список инстансов (общий для чтения и записи)."""
if not isinstance(data, list):
raise ValueError("services hosts.json: нужен список инстансов")
seen = set()
out = []
for i, ent in enumerate(data):
where = "services hosts.json запись #%d" % (i + 1)
if not isinstance(ent, dict):
raise ValueError("%s: не объект" % where)
typ = str(ent.get("type") or "").strip()
if typ not in HOSTS_TYPES:
raise ValueError("%s: type нужен один из %s" % (
where, ", ".join(HOSTS_TYPES)))
alias = str(ent.get("alias") or "").strip()
if not _ALIAS_RE.match(alias):
raise ValueError(
"%s: alias — [A-Za-z0-9_.-], до 64 символов" % where)
if alias in seen:
raise ValueError(
"%s: alias %s дублируется (должен быть уникален)" % (
where, alias))
seen.add(alias)
services = ent.get("services") or []
if not isinstance(services, list) or \
any(not str(s or "").strip() for s in services):
raise ValueError("%s: services — список непустых строк"
% where)
endpoint = str(ent.get("endpoint") or "").strip()
autostart = ent.get("autostart") or []
if not isinstance(autostart, list):
raise ValueError("%s: autostart — список" % where)
out.append({"type": typ, "alias": alias, "endpoint": endpoint,
"services": [str(s).strip() for s in services],
"autostart": [str(s).strip() for s in autostart],
"sudo_pass": str(ent.get("sudo_pass") or "")})
return out
def _read_hosts_all(self):
"""Весь hosts.json (валидированный) или [] если файла нет."""
path = self._hosts_path()
if not os.path.isfile(path):
return []
try:
with open(path, encoding="utf-8") as f:
data = json.load(f)
except Exception as e: # noqa: BLE001
raise ValueError("services hosts.json: не читается: %s" % e)
return self._validate_hosts(data)
def api_hosts_for(self, svc_type):
"""Инстансы типа из hosts.json (валидированные).
Нет файла — [] (провайдер работает как раньше: один инстанс
из своего .env). Битый файл — ValueError с понятной причиной
(setup провайдера падает с ней в лог/--check).
"""
return [e for e in self._read_hosts_all()
if e["type"] == (svc_type or "").strip()]
def api_hosts_all(self):
"""Весь hosts.json (валидированный) — для кабинета каталога."""
return self._read_hosts_all()
def api_hosts_set(self, entries, actor=""):
"""Перезаписать hosts.json целиком (атомарно) + reload провайдеров.
Кабинет: читает hosts_all, правит списки services (чекбоксы),
пишет обратно. sudo_pass/endpoint при round-trip сохраняются.
"""
entries = self._validate_hosts(entries)
path = self._hosts_path()
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(entries, f, ensure_ascii=False, indent=2)
os.replace(tmp, path)
logger.info("services: hosts.json перезаписан (%d инстансов)",
len(entries))
try:
self.ctx.api.call("audit.record", actor or "?", "services.hosts",
"%d инстансов" % len(entries), "", True)
except Exception: # noqa: BLE001
pass
reloaded = []
for provider in self._providers:
try:
self.ctx.api.call(provider + ".reload")
reloaded.append(provider)
except Exception as e: # noqa: BLE001
logger.warning("services: reload %s: %s", provider, e)
return {"ok": True, "hosts": len(entries), "reloaded": reloaded}
# --- кеш каталога (data/services/discovery_cache.json) ---
def _load_cache(self):
path = self._cache_path()
if not os.path.isfile(path):
return {}
try:
with open(path, encoding="utf-8") as f:
data = json.load(f)
except Exception: # noqa: BLE001
return {}
return data if isinstance(data, dict) else {}
def _save_cache(self, cache):
path = self._cache_path()
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(cache, f, ensure_ascii=False, indent=1)
os.replace(tmp, path)
def _host_map(self):
"""{alias: (provider, type)} по живым инстансам провайдеров."""
out = {}
for provider in self._providers:
for info in self._infos(provider):
host = info.get("host") or ""
if host and host not in out:
out[host] = (provider, info.get("type") or "?")
return out
def api_discover(self, host="", actor=""):
"""Опросить хосты и перезаписать кеш каталога (только ручной запуск).
host пустой — все хосты всех провайдеров. Возврат
{updated: [...], errors: {host: err}}. Ошибки одного хоста не
валят остальных; его старый кеш сохраняется.
"""
host = (host or "").strip()
targets = self._host_map()
if host:
if host not in targets:
raise ValueError("хост %s не найден" % host)
targets = {host: targets[host]}
if not targets:
return {"updated": [], "errors": {}}
cache = self._load_cache()
now = int(_time.time())
updated, errors = [], {}
for alias, (provider, typ) in sorted(targets.items()):
try:
found = self.ctx.api.call(provider + ".discover", alias)
except Exception as e: # noqa: BLE001
errors[alias] = "%s: %s" % (type(e).__name__, e)
continue
names = []
for one in found if isinstance(found, list) else []:
if isinstance(one, dict) and \
(one.get("host") or "") == alias:
if one.get("error"):
errors[alias] = str(one["error"])[:300]
# старый кеш не затираем пустым — только штамп ошибки
prev = (cache.get("%s:%s" % (typ, alias)) or {})
names = prev.get("services") or []
else:
names = one.get("services") or []
cache["%s:%s" % (typ, alias)] = {
"at": now, "host": alias, "type": typ,
"services": [str(s) for s in names if str(s or "")],
"error": errors.get(alias, "")}
updated.append(alias)
self._save_cache(cache)
try:
self.ctx.api.call("audit.record", actor or "?",
"services.discover",
"hosts=%s" % ",".join(sorted(targets)),
"ok=%d err=%d" % (len(updated),
len(errors)),
not errors)
except Exception: # noqa: BLE001
pass
return {"updated": updated, "errors": errors}
def api_catalog(self, query="", type_="", host="", checked="all",
page=1, per_page=CATALOG_DEFAULT_PER_PAGE):
"""Каталог из кеша: фильтры + пагинация.
checked: all|on|off (в вайтлисте инстанса или нет).
Возврат {items: [{host,type,name,checked,at}], total, page,
per_page, pending: [{host,type}] (хосты без кеша), cache_at}.
"""
q = (query or "").strip().lower()
tf = (type_ or "").strip()
hf = (host or "").strip()
cf = (checked or "all").strip()
if cf not in ("all", "on", "off"):
raise ValueError("checked: all|on|off")
try:
page = 1 if page is None else int(page)
except (TypeError, ValueError):
raise ValueError("page — число >= 1")
if page < 1:
raise ValueError("page — число >= 1")
try:
per_page = int(per_page or CATALOG_DEFAULT_PER_PAGE)
except (TypeError, ValueError):
raise ValueError("per_page — число")
per_page = min(max(1, per_page), CATALOG_MAX_PER_PAGE)
checked_sets = {}
order = []
for provider in self._providers:
for info in self._infos(provider):
h = info.get("host") or ""
if h and h not in checked_sets:
checked_sets[h] = set(info.get("services") or [])
order.append((info.get("type") or "?", h))
cache = self._load_cache()
type_rank = {t: i for i, t in enumerate(HOSTS_TYPES)}
items = []
cached_hosts = set()
for key, ent in cache.items():
if not isinstance(ent, dict):
continue
h = str(ent.get("host") or "")
if not h or h not in checked_sets:
continue # протухший кеш удалённого хоста
cached_hosts.add(h)
t = str(ent.get("type") or "?")
at = ent.get("at")
in_list = checked_sets[h]
for name in ent.get("services") or []:
name = str(name or "")
if not name:
continue
is_on = name in in_list
if q and q not in name.lower():
continue
if tf and t != tf:
continue
if hf and h != hf:
continue
if cf == "on" and not is_on:
continue
if cf == "off" and is_on:
continue
items.append({"host": h, "type": t, "name": name,
"checked": is_on, "at": at})
items.sort(key=lambda r: (type_rank.get(r["type"], 99),
r["host"], r["name"].lower()))
total = len(items)
start = (page - 1) * per_page
pending = [{"host": h, "type": t} for t, h in order
if h not in cached_hosts]
cache_at = {}
for key, ent in cache.items():
if isinstance(ent, dict) and (ent.get("host") or "") in \
checked_sets:
cache_at[ent["host"]] = ent.get("at")
return {"items": items[start:start + per_page], "total": total,
"page": page, "per_page": per_page, "pending": pending,
"cache_at": cache_at}
# --- точки outbound (создание из каталога) ---
#: Тип инстанса -> kind точки (local без точки — команды локальные).
ENDPOINT_KINDS = {"windows": "ssh", "linux": "ssh", "docker": "ssh"}
def api_endpoint_add(self, name, svc_type, address, secret="",
via="", note="", actor=""):
"""Создать точку outbound для будущего инстанса + привязать выбор.
Секрет (пароль SSH) — через modsafe.env_set в data/env/outbound.env
под auth_ref вида SVC_EP_<NAME> (в логах маскируется). Проба
подключения — внутри outbound.endpoint_add (точка не отвечает —
ValueError, ничего не сохраняется). Возврат {ok, name}.
"""
ep_name = (name or "").strip()
svc_type = (svc_type or "").strip()
if svc_type not in self.ENDPOINT_KINDS:
raise ValueError("svc_type: windows|linux|docker "
"(local — без точки)")
if not isinstance(address, dict):
raise ValueError("address — объект {host, port?, user}")
kind = self.ENDPOINT_KINDS[svc_type]
auth_ref = ""
if (secret or "").strip():
auth_ref = "SVC_EP_%s" % _re.sub(
r"[^A-Z0-9]", "_", ep_name.upper())[:48]
try:
self.ctx.api.call("modsafe.env_set", "outbound",
{auth_ref: secret}, actor or "?")
except Exception as e: # noqa: BLE001
raise ValueError("не saved секрет точки: %s" % e)
try:
res = self.ctx.api.call(
"outbound.endpoint_add", ep_name, kind, dict(address),
**({"auth_ref": auth_ref} if auth_ref else {}),
**({"via": (via or "").strip()} if (via or "").strip()
else {}),
**({"note": (note or "")[:500]} if (note or "").strip()
else {}))
except Exception as e: # noqa: BLE001
raise ValueError(str(e)[:300])
if not isinstance(res, dict) or not res.get("ok"):
err = (res or {}).get("error") if isinstance(res, dict) \
else "bad result"
raise ValueError(str(err)[:300])
try:
self.ctx.api.call("audit.record", actor or "?",
"services.endpoint", ep_name,
"kind=%s type=%s" % (kind, svc_type), True)
except Exception: # noqa: BLE001
pass
return {"ok": True, "name": ep_name}
def api_endpoint_delete(self, name, actor=""):
"""Удалить точку outbound + сбросить ссылки в hosts.json.
Точка может использоваться инстансами (поле endpoint): их ссылки
сбрасываются в "" (с reload провайдеров) — иначе хосты остались
бы с битой привязкой. Секрет SVC_EP_<NAME> в outbound.env
остаётся сиротой (безвреден, маскируется; чистится в кабинете
outbound/env при желании). Возврат {ok, name, cleared: [alias]}.
"""
ep_name = (name or "").strip()
if not ep_name:
raise ValueError("пустое имя точки")
entries = self._read_hosts_all()
cleared = [e["alias"] for e in entries if e.get("endpoint")
== ep_name]
if cleared:
for e in entries:
if e.get("endpoint") == ep_name:
e["endpoint"] = ""
self.api_hosts_set(entries, actor)
try:
res = self.ctx.api.call("outbound.endpoint_delete", ep_name)
except Exception as e: # noqa: BLE001
raise ValueError(str(e)[:300])
if isinstance(res, dict) and res.get("ok") is False:
raise ValueError(str(res.get("error") or "?")[:300])
try:
self.ctx.api.call("audit.record", actor or "?",
"services.endpoint_delete", ep_name,
"cleared=%s" % ",".join(cleared), True)
except Exception: # noqa: BLE001
pass
return {"ok": True, "name": ep_name, "cleared": cleared}
def api_host_add(self, alias, svc_type, endpoint="", sudo_pass="",
actor=""):
"""Добавить инстанс-хост в hosts.json (для local — без точки).
local: endpoint обязан быть пустым (команды локальные, outbound
не нужен). windows/linux/docker: endpoint может быть пустым
(привяжете позже селектом в таблице) или именем существующей
точки. Новый хост стартует с пустым services — дальше discovery
покажет каталог, чекбоксы наберут вайтлист. Возврат {ok, alias}.
"""
alias = (alias or "").strip()
svc_type = (svc_type or "").strip()
endpoint = (endpoint or "").strip()
if not _ALIAS_RE.match(alias or ""):
raise ValueError("alias — [A-Za-z0-9_.-], до 64 символов")
if svc_type not in HOSTS_TYPES:
raise ValueError("svc_type: %s" % ", ".join(HOSTS_TYPES))
if svc_type == "local" and endpoint:
raise ValueError("local — без точки (endpoint пустой)")
entries = self._read_hosts_all()
if any(e["alias"] == alias for e in entries):
raise ValueError("alias %s уже есть" % alias)
if endpoint:
try:
eps = self.ctx.api.call("outbound.endpoint_list") or []
except Exception as e: # noqa: BLE001
raise ValueError("нет списка точек outbound: %s" % e)
names = {e.get("name") for e in eps
if isinstance(e, dict)}
if endpoint not in names:
raise ValueError("точки %s нет (создайте её выше)"
% endpoint)
entries.append({"type": svc_type, "alias": alias,
"endpoint": endpoint, "services": [],
"autostart": [],
"sudo_pass": (sudo_pass or "")})
self.api_hosts_set(entries, actor)
return {"ok": True, "alias": alias}
# --- реестр ---
def api_register(self, provider):
"""Регистрация провайдера (имя модуля). Идемпотентно."""
provider = (provider or "").strip()
if not provider:
raise ValueError("register: пустое имя провайдера")
self._providers.setdefault(provider, True)
logger.info("services: провайдер %s", provider)
return True
def _infos(self, provider):
"""Список инстансов провайдера (dict от старого провайдера — нормализуем)."""
try:
info = self.ctx.api.call(provider + ".services_info")
except Exception as e: # noqa: BLE001
return [{"host": provider, "type": "?", "services": [],
"error": "%s: %s" % (type(e).__name__, e)}]
if isinstance(info, dict):
info = [info]
if not isinstance(info, list):
return [{"host": provider, "type": "?", "services": [],
"error": "bad info"}]
out = []
for one in info:
if not isinstance(one, dict):
out.append({"host": provider, "type": "?", "services": [],
"error": "bad info"})
continue
one = dict(one)
one.setdefault("host", provider)
one.setdefault("type", "?")
one.setdefault("services", [])
out.append(one)
return out
def _info(self, provider):
"""Совместимость: первый инстанс (для одиночных ответов)."""
infos = self._infos(provider)
return infos[0] if infos else {"host": provider, "type": "?",
"services": []}
def api_providers(self):
"""{модуль: [info...]} всех зарегистрированных провайдеров."""
return {p: self._infos(p) for p in self._providers}
# --- агрегаты ---
def api_list(self):
"""Плоский список [{host, name, type}] + просто имена."""
items = []
for provider in self._providers:
for info in self._infos(provider):
for name in info.get("services") or []:
items.append({"host": info.get("host"), "name": name,
"type": info.get("type")})
items.sort(key=lambda x: (x["host"] or "", x["name"] or ""))
return {"services": items,
"names": sorted({i["name"] for i in items if i["name"]})}
def api_status(self):
"""Статусы всех служб всех провайдеров (best-effort)."""
items = []
for provider in self._providers:
for info in self._infos(provider):
host = info.get("host")
for name in info.get("services") or []:
try:
st = self.ctx.api.call(provider + ".service_status",
name, host)
except Exception as e: # noqa: BLE001
st = {"name": name, "up": False,
"detail": "%s: %s" % (type(e).__name__, e)}
if not isinstance(st, dict):
st = {"name": name, "up": False,
"detail": "bad status"}
items.append({"host": host, "name": name,
"type": info.get("type"),
"status": bool(st.get("up")),
"detail": str(st.get("detail") or "")})
return {"services": items}
def api_action(self, host, name, act, actor=""):
"""Действие start|stop|restart: маршрутизация по host."""
act = (act or "").strip()
if act not in ("start", "stop", "restart"):
raise ValueError("недопустимое действие (start|stop|restart)")
for provider in self._providers:
for info in self._infos(provider):
if (info.get("host") or "") == (host or "") and \
name in (info.get("services") or []):
try:
res = self.ctx.api.call(provider + ".service_action",
name, act, host)
except Exception as e: # noqa: BLE001
res = {"ok": False,
"error": "%s: %s" % (type(e).__name__, e)}
if not isinstance(res, dict):
res = {"ok": False, "error": "bad action result"}
try:
self.ctx.api.call(
"audit.record", actor or "?",
"service.action", "%s/%s" % (host, name),
"%s ok=%s" % (act, bool(res.get("ok"))),
bool(res.get("ok")))
except Exception: # noqa: BLE001
pass
return res
raise ValueError("служба %s/%s не найдена" % (host, name))
# --- встроенные провайдеры (бывшие модули services_*) ---
try:
from .backends.provider import ProviderModule as _ProviderModule
from .backends.docker import DockerServices as _DockerServices
from .backends.linux import LinuxServices as _LinuxServices
from .backends.local import LocalServices as _LocalServices
from .backends.windows import WinServices as _WinServices
except ImportError:
# Прямой setup без пакета (тесты import modules.services.module):
# backends подтянется через modules.services.backends.
from modules.services.backends.provider import ( # noqa: F401
ProviderModule as _ProviderModule)
from modules.services.backends.docker import ( # noqa: F401
DockerServices as _DockerServices)
from modules.services.backends.linux import ( # noqa: F401
LinuxServices as _LinuxServices)
from modules.services.backends.local import ( # noqa: F401
LocalServices as _LocalServices)
from modules.services.backends.windows import ( # noqa: F401
WinServices as _WinServices)
class _EmbeddedProvider(_ProviderModule):
"""Провайдер внутри ядра: setup вызывается ядром, не loader'ом."""
name = "services-provider"
def __init__(self, mod_name, mod_dir, ctx):
self._mod_name = mod_name
self._mod_dir = mod_dir
self.setup(ctx)
class _WindowsProvider(_EmbeddedProvider):
MOD = "services_windows"
TYPE = "windows"
ENGINE = _WinServices
ENGINE_ARGS = "services+autostart"
NEED_OUTBOUND = True
CONSUMER_DESC = "управление службами Windows (sc)"
LEGACY_ALIAS_KEY = "WIN_ALIAS"
LEGACY_ALIAS_DEFAULT = "win"
LEGACY_SERVICES_KEY = "WIN_SERVICES"
LEGACY_EXTRA = {"autostart": "WIN_AUTOSTART"}
UNIT_NOUN = "сервис"
HEALTH_COUNT = "services"
class _LinuxProvider(_EmbeddedProvider):
MOD = "services_linux"
TYPE = "linux"
ENGINE = _LinuxServices
ENGINE_ARGS = "services+sudo"
NEED_OUTBOUND = True
CONSUMER_DESC = "управление юнитами Linux (systemctl)"
LEGACY_ALIAS_KEY = "LINUX_ALIAS"
LEGACY_ALIAS_DEFAULT = "linux"
LEGACY_SERVICES_KEY = "LINUX_SERVICES"
LEGACY_EXTRA = {"sudo_pass": "LINUX_SUDO_PASS"}
UNIT_NOUN = "юнит"
HEALTH_COUNT = "units"
class _DockerProvider(_EmbeddedProvider):
MOD = "services_docker"
TYPE = "docker"
ENGINE = _DockerServices
ENGINE_ARGS = "containers"
NEED_OUTBOUND = True
CONSUMER_DESC = "управление контейнерами Docker"
LEGACY_ALIAS_KEY = "DOCKER_ALIAS"
LEGACY_ALIAS_DEFAULT = "docker"
LEGACY_SERVICES_KEY = "DOCKER_CONTAINERS"
UNIT_NOUN = "контейнер"
HEALTH_COUNT = "containers"
class _LocalProvider(_EmbeddedProvider):
MOD = "services_local"
TYPE = "local"
ENGINE = _LocalServices
ENGINE_ARGS = "services+sudo-local"
NEED_OUTBOUND = False
LEGACY_ALIAS_KEY = "LOCAL_ALIAS"
LEGACY_ALIAS_DEFAULT = "local"
LEGACY_SERVICES_KEY = "LOCAL_SERVICES"
LEGACY_EXTRA = {"sudo_pass": "LOCAL_SUDO_PASS"}
UNIT_NOUN = "юнит"
HEALTH_COUNT = "units"
_EMBEDDED_CLASSES = {
"services_windows": _WindowsProvider,
"services_linux": _LinuxProvider,
"services_docker": _DockerProvider,
"services_local": _LocalProvider,
}
tg.py 6.2 КБ
"""services.tg — Telegram-меню служб (группировка по хостам).
Данные — через ctx.api ядра services (реестр провайдеров). Действия —
в отдельном потоке (SSH/systemctl могут висеть), результат — новым сообщением.
"""
import logging
import threading
logger = logging.getLogger(__name__)
class ServiceMenu:
def __init__(self, ctx):
self.ctx = ctx
# --- helpers ---
def _can(self, user_id):
try:
return bool(self.ctx.rights.can("telegram:%s" % user_id,
"service_manage"))
except Exception: # noqa: BLE001
return False
def _send(self, bot, chat_id, text, reply_markup=None):
from telegram import InlineKeyboardMarkup # noqa: F401
try:
bot.send_message(chat_id=chat_id, text=text,
reply_markup=reply_markup,
disable_web_page_preview=True)
except Exception as e: # noqa: BLE001
logger.warning("services tg send: %s", e)
# --- экраны ---
def show_menu(self, update, context):
from telegram import InlineKeyboardButton, InlineKeyboardMarkup
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
try:
data = self.ctx.api.call("services.status")
items = (data or {}).get("services") or []
except Exception as e: # noqa: BLE001
self._send(context.bot, chat_id, "Ошибка статуса служб: %s" % e)
return
if not items:
self._send(context.bot, chat_id, "Нет служб, доступных для управления.")
return
keyboard = []
row = []
for it in items:
emoji = "🟢" if it.get("status") else "🔴"
title = "%s %s/%s" % (emoji, it.get("host"), it.get("name"))
row.append(InlineKeyboardButton(
title, callback_data="service_details:%s:%s" % (
it.get("host"), it.get("name"))))
if len(row) == 1:
keyboard.append(row)
row = []
if row:
keyboard.append(row)
keyboard.append([InlineKeyboardButton("🔙 Назад", callback_data="main_menu")])
self._send(context.bot, chat_id, "🔩 Службы (хост/служба):",
InlineKeyboardMarkup(keyboard))
def show_details(self, update, context, host, name):
from telegram import InlineKeyboardButton, InlineKeyboardMarkup
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
kb = [
[InlineKeyboardButton(
"🟢 Запустить",
callback_data="service_action:%s:%s:start" % (host, name)),
InlineKeyboardButton(
"⛔ Выключить",
callback_data="service_action:%s:%s:stop" % (host, name))],
[InlineKeyboardButton(
"♻️ Перезапустить",
callback_data="service_action:%s:%s:restart" % (host, name))],
[InlineKeyboardButton("🔙 Назад", callback_data="service_menu")],
]
self._send(context.bot, chat_id, "🔩 %s/%s — действие:" % (host, name),
InlineKeyboardMarkup(kb))
def do_action(self, update, context, host, name, act):
from telegram.ext import DispatcherHandlerStop
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can(user_id):
try:
update.callback_query.edit_message_text("🚫 Доступ запрещён!")
except Exception: # noqa: BLE001
pass
raise DispatcherHandlerStop()
self._send(context.bot, chat_id,
"Команда отправлена на выполнение. Ожидайте результата...")
ctx = self.ctx
def _run():
try:
res = ctx.api.call("services.action", host, name, act,
"telegram:%s" % user_id)
except Exception as e: # noqa: BLE001
res = {"ok": False, "error": "%s: %s" % (type(e).__name__, e)}
if res.get("ok"):
text = "Результат %s для %s/%s:\n%s" % (
act, host, name, res.get("output") or "OK")
else:
text = "Ошибка %s для %s/%s: %s" % (
act, host, name, res.get("error") or "?")
try:
ctx.api.call("transport_telegram.send", chat_id, text)
except Exception: # noqa: BLE001
try:
ctx.api.call("notify.send", "telegram:%s" % user_id, text)
except Exception: # noqa: BLE001
logger.warning("services tg: некуда ответить %s", chat_id)
threading.Thread(target=_run, daemon=True,
name="svc-action").start()
raise DispatcherHandlerStop()
# --- вход ---
def handle(self, update, context):
from telegram.ext import DispatcherHandlerStop
try:
data = update.callback_query.data or ""
update.callback_query.answer()
except Exception: # noqa: BLE001
return
if data == "service_menu":
self.show_menu(update, context)
elif data.startswith("service_details:"):
parts = data.split(":", 2)
if len(parts) == 3:
self.show_details(update, context, parts[1], parts[2])
elif data.startswith("service_action:"):
parts = data.split(":", 3)
if len(parts) == 4:
self.do_action(update, context, parts[1], parts[2], parts[3])
return
raise DispatcherHandlerStop()
ui_tg.py 0.8 КБ
"""services.ui_tg — кнопка /start + меню служб (колбэки service_*).
Группа -1: catch-all admin.button (группа 0) не должен съедать наши колбэки.
"""
TG_COMMANDS = {}
TG_MENU = [
{"id": "services", "title": "🔩 Сервисы", "callback": "service_menu",
"right": "service_manage"},
]
def register_ptb(facade, mod):
"""Меню служб на фасаде напрямую (колбэки service_*), группа -1."""
from telegram.ext import CallbackQueryHandler
from .tg import ServiceMenu
menu = ServiceMenu(mod.ctx)
mod.svc_menu = menu
facade.add_handler(CallbackQueryHandler(
menu.handle,
pattern=r"^(service_menu|service_details:|service_action:)"), group=-1)
ui_web.py 11.4 КБ
"""services.ui_web — управление службами через ядро (реестр провайдеров).
GET /api/admin/services — статусы [{host, name, type, status, detail}].
POST /api/admin/service — start/stop/restart {host?, name, action}.
host опционален, если имя уникально; при дублях — 400 с подсказкой.
Каталог (право service_catalog): GET /api/admin/services_catalog (фильтры +
пагинация, данные из кеша discovery), POST /api/admin/services_discover
{host?} (ручное обновление кеша, в фоне), GET /api/admin/services_hosts
(инстансы + привязки), POST /api/admin/services_hosts {entries}
(чекбоксы видимости, без рестарта), POST /api/admin/services_endpoint
{cоздание точки outbound для инстанса}.
"""
ROUTES = [
("GET", "/api/admin/services", "services", {"right": "service_manage",
"desc": "Статусы служб всех провайдеров"}),
("POST", "/api/admin/service", "service", {"right": "service_manage",
"desc": "start/stop/restart службы {action, host?, name} (host обязателен при дублях имени)"}),
("GET", "/api/admin/services_catalog", "catalog",
{"right": "service_catalog",
"desc": "Каталог служб из кеша: фильтры + пагинация"}),
("POST", "/api/admin/services_discover", "discover",
{"right": "service_catalog",
"desc": "Обновить кеш каталога: {host?}, в фоне"}),
("GET", "/api/admin/services_hosts", "hosts",
{"right": "service_catalog",
"desc": "Инстансы hosts.json + привязки outbound"}),
("POST", "/api/admin/services_hosts", "hosts_save",
{"right": "service_catalog",
"desc": "Сохранить чекбоксы видимости: {entries}"}),
("POST", "/api/admin/services_endpoint", "endpoint",
{"right": "service_catalog",
"desc": "Создать точку outbound: {name, svc_type, address, secret?, via?, note?}"}),
("POST", "/api/admin/services_endpoint_delete", "endpoint_delete",
{"right": "service_catalog",
"desc": "Удалить точку outbound: {name} (ссылки хостов сбрасываются)"}),
("POST", "/api/admin/services_host_add", "host_add",
{"right": "service_catalog",
"desc": "Добавить хост-инстанс: {alias, svc_type, endpoint?, sudo_pass?}"}),
]
SERVICE = {"key": "service_manage", "icon": "🔧", "title": "Управление службами", "right": "service_manage",
"phase": 3, "desc": "Службы хостов (Windows/Docker/локально): старт/стоп/рестарт"}
PANELS = []
def _need_catalog(ctx, req):
"""403 если у web-юзера нет права service_catalog (окно каталога)."""
sess = req.get("session") or {}
uid = sess.get("uid")
try:
if uid is not None and ctx.rights.can("web:%s" % uid,
"service_catalog"):
return None
except Exception: # noqa: BLE001
pass
return 403, {"ok": False, "error": "Нужно право service_catalog"}
def handle_api(ctx, config, method, req):
from core.errors import UserError
if method == "services":
try:
return ctx.api.call("services.status")
except Exception as e: # noqa: BLE001
return 500, {"error": "%s: %s" % (type(e).__name__, e)}
if method == "service":
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
name = str(body.get("name") or "").strip()
host = str(body.get("host") or "").strip()
act = str(body.get("action") or "").strip()
if act not in ("start", "stop", "restart"):
return 400, {"ok": False, "error": "Недопустимое действие"}
if not name:
return 400, {"ok": False, "error": "Нет имени службы"}
if not host:
# однозначное имя — резолвим сами; дубли требуют host
try:
items = (ctx.api.call("services.list") or {}).get("services") or []
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": str(e)[:200]}
hosts = sorted({i.get("host") for i in items
if i.get("name") == name})
if not hosts:
return 400, {"ok": False, "error": "Служба не найдена"}
if len(hosts) > 1:
return 400, {"ok": False,
"error": "Имя есть на нескольких хостах: %s (укажите host)"
% ", ".join(hosts)}
host = hosts[0]
try:
res = ctx.api.call("services.action", host, name, act, actor)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (type(e).__name__, e)}
if not isinstance(res, dict) or not res.get("ok"):
err = (res or {}).get("error") if isinstance(res, dict) else None
return 500, {"ok": False, "error": str(err or "?")[:600]}
return {"ok": True, "output": str(res.get("output") or "OK")[:2000]}
if method == "catalog":
denied = _need_catalog(ctx, req)
if denied:
return denied
query = req.get("query") or {}
try:
return ctx.api.call(
"services.catalog",
query.get("q") or "", query.get("type") or "",
query.get("host") or "", query.get("checked") or "all",
query.get("page") or 1, query.get("per_page") or 50)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (
type(e).__name__, e)}
if method == "discover":
denied = _need_catalog(ctx, req)
if denied:
return denied
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
host = str(body.get("host") or "").strip()
import threading as _th
def _run():
try:
ctx.api.call("services.discover", host, actor)
except Exception: # noqa: BLE001
pass
_th.Thread(target=_run, daemon=True,
name="svc-discover").start()
return {"ok": True, "started": True, "host": host or "all"}
if method == "hosts":
denied = _need_catalog(ctx, req)
if denied:
return denied
try:
entries = ctx.api.call("services.hosts_all")
except ValueError as e:
return 500, {"ok": False, "error": str(e)}
bound = {}
try:
for c in ctx.api.call("outbound.consumer_list") or []:
if isinstance(c, dict) and str(c.get("module") or "")\
.startswith("services_"):
bound[c["module"]] = c.get("bound")
except Exception: # noqa: BLE001
pass
endpoints = []
try:
for e in ctx.api.call("outbound.endpoint_list") or []:
if isinstance(e, dict) and e.get("name"):
endpoints.append({"name": e["name"],
"kind": e.get("kind") or ""})
except Exception: # noqa: BLE001
pass
return {"ok": True, "hosts": entries, "bindings": bound,
"endpoints": sorted(endpoints, key=lambda e: e["name"])}
if method == "hosts_save":
denied = _need_catalog(ctx, req)
if denied:
return denied
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
entries = body.get("entries")
if not isinstance(entries, list):
return 400, {"ok": False, "error": "entries — список инстансов"}
try:
return ctx.api.call("services.hosts_set", entries, actor)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (
type(e).__name__, e)}
if method == "endpoint":
denied = _need_catalog(ctx, req)
if denied:
return denied
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
address = body.get("address")
if not isinstance(address, dict):
# плоская форма из кабинета — собираем address сами
address = {}
for k in ("host", "port", "user", "key"):
v = str(body.get(k) or "").strip()
if v:
address[k] = v
try:
return ctx.api.call(
"services.endpoint_add",
str(body.get("name") or "").strip(),
str(body.get("svc_type") or "").strip(),
address,
str(body.get("secret") or ""),
str(body.get("via") or ""),
str(body.get("note") or ""), actor)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (
type(e).__name__, e)}
if method == "endpoint_delete":
denied = _need_catalog(ctx, req)
if denied:
return denied
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
name = str(body.get("name") or "").strip()
if not name:
return 400, {"ok": False, "error": "Нет имени точки"}
try:
return ctx.api.call("services.endpoint_delete", name, actor)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (
type(e).__name__, e)}
if method == "host_add":
denied = _need_catalog(ctx, req)
if denied:
return denied
body = req.get("body") or {}
sess = req.get("session") or {}
actor = "web:%s" % sess.get("uid") if sess.get("uid") else "?"
try:
return ctx.api.call(
"services.host_add",
str(body.get("alias") or "").strip(),
str(body.get("svc_type") or "").strip(),
str(body.get("endpoint") or "").strip(),
str(body.get("sudo_pass") or ""), actor)
except ValueError as e:
return 400, {"ok": False, "error": str(e)}
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": "%s: %s" % (
type(e).__name__, e)}
raise UserError("неизвестный метод: %s" % method)