network v0.3.8 [extra]
Сеть: сканер устройств, WireGuard/Guacamole/fail2ban-мониторы, MAC-блокировки, сниффер.
| version | date | commit | файлов |
|---|---|---|---|
| 0.3.8 | 2026-09-16 | 2d7d04615b31 | 13 |
README
# network `0.3.8`
Карта домашней сети: кто онлайн, кто новый, кто пропал.
Тип: Модуль. Категория: `extra`. Зависимостей нет.
## Описание
Глаза бота в локалке `10.10.10.0/24`. По кругу ходит nmap, сверяется с кешем и докладывает: новое устройство, переподключение, отключение (только после нескольких пропусков подряд — спасает дебаунс), смена IP. Мелкие повторы склеиваются агрегацией. Каждому устройству копятся статистика и открытые порты, имена подписываются вручную (`nmap_setname` принимается и из TG, и из веба), знакомые MAC — в игнор-лист.
Доп. источники (fail2ban, WireGuard, Guacamole — код в своих модулях) подключаются через `register_monitor`. Длинные списки в TG режутся на чанки по 3000 символов.
## Возможности
- `scanner` `{}` → движок сканирования (кеш устройств)
- `cache_file` `{}` → путь кеша устройств
- `ignored_file` `{}` → путь игнор-листа
- `devices_snapshot` `{}` → `{total, online, devices[]}` — обогащённый снимок для мониторинга
- `register_monitor` `{name, factory, interval_s, type, silent}` → плагин доп. источника
## Права
- `nmap_manage` — Сканер сети (GET/POST /api/admin/nmap*)
- `receive_notifications_nmap` — Уведомления сканера сети (новые/потерянные устройства)
- `receive_notifications_nmap_edit` — Право менять себе уведомления сканера сети
## Telegram
- кнопка «📡 Сканирование» → меню, постраничные списки устройств и статистики (просмотр — право `view_net_status`, изменения — `nmap_manage`)
- переименование устройств диалогом (`net_setname` / `nmap_setname`)
## Web
- `GET /api/network/devices` — устройства для мониторинга
- `GET /api/admin/nmap` — устройства + игнор-лист
- `POST /api/admin/nmap/name` — подписать имя `{mac, name}`
- `POST /api/admin/nmap/ignore` — игнор `{action: add|remove, mac}`
- `POST /api/network/nmap/delete` — удалить устройство (только веб, в TG удаления нет)
- карточка «Сканер сети»
## События шины
- излучает `network.device_new`, `network.device_reconnected`, `network.device_disconnected`, `network.device_ip_changed` (по одному на устройство)
## Конфигурация (`config.schema.json`)
- обязательных нет; опциональные (дефолты):
- `SCAN_NETWORK_RANGE` = `10.10.10.0/24` — диапазон сканирования
- `CHECK_DURATION_NSCAN_LONG` = `30` — длинный цикл, сек
- `CHECK_DURATION_NSCAN_FAST` = `5` — быстрый цикл, сек
- `DEBOUNCE_COUNT` = `2` — пропусков до «пропал»
- `AGGREGATION_INTERVAL` = `300` — окно склейки уведомлений, сек
- `AGGREGATION_MAX_CHARS` = `3500` — лимит склеенного текста
- `NOTIF_MAX_RETRIES` = `3` — ретраи уведомлений
- `NOTIF_RETRY_BACKOFF` = `1.5` — backoff ретраев
- `F2B_ENABLED` = `1` — штатный монитор fail2ban вкл
Манифест
{
"name": "network",
"version": "0.3.8",
"category": "extra",
"description": "Сеть: сканер устройств, WireGuard/Guacamole/fail2ban-мониторы, MAC-блокировки, сниффер.",
"requires": [],
"rights": [],
"api_methods": [],
"web_routes": [],
"commit": "2d7d04615b31be9f5bbd91d5d61671182164bf39",
"updated": "2026-09-16 16:35:50 +0300",
"integrity": "sha256:d3b1ccc1fee0850c875620129b6dde2063fa3aa066fbc98ff4d793ce3938a6ef",
"tree_url": "http://10.10.10.30:3000/euk0r/bot_modular/src/branch/main/modules/network",
"raw_url": "http://10.10.10.30:3000/euk0r/bot_modular/raw/branch/main/modules/network/module.py",
"files": "[13 файлов — см. вкладки ниже]",
"note": "Модуль проверки локальной сети",
"_extra": {
"api": "",
"events_emitted": "",
"events_subscribed": "[]",
"audit_allow": "",
"uses_rights": "",
"shared_routes": "",
"config_schema": "config.schema.json",
"env_example": ".env.example"
}
} Файлы и исходники
Дерево файлов
- · корень
- 3.8 КБ
- 2.8 КБ
- 0.3 КБ
- 6.6 КБ
- 2.3 КБ
- 11.3 КБ
- 38.4 КБ
- 6.1 КБ
- 16.3 КБ
- 1.2 КБ
- 4.5 КБ
- web/
- 1.1 КБ
- 6.0 КБ
Предпросмотр
Выберите файл в дереве выше — код откроется здесь.
Все исходники (.py) одним списком
config.py 2.8 КБ
"""Пакетный config-шим для network (scan пишет config.X).
ВНИМАНИЕ: macblock импортирует значения (from .config import Y) — поэтому
install() обязан вызываться ДО импорта движков (порядок в module.py).
Пути кешей — data/cache/network (свежие, пишет сам сканер); зеркало
прод-кеша для чтения — через MONITOR_CACHE_DIR транспорта (только чтение).
Цели уведомлений сканера — единый ADMIN_IDS модуля access (через ctx.api).
"""
import os
_vals: dict = {}
def _int(values: dict, key: str, default: int) -> int:
try:
return int(str(values.get(key, default) or default))
except (ValueError, TypeError):
return default
def install(values: dict, files: dict) -> None:
_vals.clear()
_vals.update(values)
_vals.update(files)
for k, d in (("CHECK_DURATION_NSCAN_LONG", 30), ("CHECK_DURATION_NSCAN_FAST", 5),
("DEBOUNCE_COUNT", 2), ("AGGREGATION_INTERVAL", 300),
("AGGREGATION_MAX_CHARS", 3500), ("NOTIF_MAX_RETRIES", 3),
("ROUTER_PORT", 22), ("TRAFFIC_MAX_SCANS", 50),
("TRAFFIC_MAX_PCAP_MB", 20), ("TRAFFIC_DEFAULT_DURATION", 30),
("TRAFFIC_MAX_DURATION", 120), ("TRAFFIC_DEFAULT_PACKETS", 2000),
("TRAFFIC_MAX_PACKETS", 10000)):
_vals[k] = _int(values, k, d)
try:
_vals["NOTIF_RETRY_BACKOFF"] = float(
str(values.get("NOTIF_RETRY_BACKOFF", "1.5") or "1.5"))
except (ValueError, TypeError):
_vals["NOTIF_RETRY_BACKOFF"] = 1.5
for k in ("ROUTER_HOST", "ROUTER_USER", "ROUTER_PASSWORD",
"SCAN_NETWORK_RANGE", "F2B_ENABLED"):
_vals.setdefault(k, (values.get(k) or "").strip())
if "ROUTER_HOST" not in _vals or not _vals["ROUTER_HOST"]:
_vals["ROUTER_HOST"] = (values.get("ROUTER_HOST") or "10.10.10.1").strip()
if not _vals.get("ROUTER_USER"):
_vals["ROUTER_USER"] = (values.get("ROUTER_USER") or "root").strip()
if not _vals.get("SCAN_NETWORK_RANGE"):
_vals["SCAN_NETWORK_RANGE"] = (
values.get("SCAN_NETWORK_RANGE") or "10.10.10.0/24").strip()
prot = (values.get("MACBLOCK_PROTECTED") or "bc:24:11:39:ea:6c").strip()
_vals["MACBLOCK_PROTECTED"] = [m.strip().upper() for m in prot.split(",") if m.strip()]
_vals["ADMIN_PASS"] = values.get("ADMIN_PASS") or ""
for k, d in (("TRAFFIC_DEFAULT_DURATION", 30), ("TRAFFIC_DEFAULT_PACKETS", 2000)):
_vals[k] = _int(values, k, d)
def __getattr__(name: str):
try:
return _vals[name]
except KeyError:
raise AttributeError("network.config: нет ключа %s" % name) from None
fail2ban_mon.py 6.6 КБ
import subprocess
import json
import os
from datetime import datetime
import logging
logger = logging.getLogger(__name__)
def parse_event_day(event_time):
"""
Извлекает дату из event_time.
Ожидаемый формат: "YYYY-MM-DD HH:MM:SS", возвращает "YYYY-MM-DD".
"""
return event_time.split()[0]
class Fail2BanMonitor:
def __init__(self, cache_file, monitoring_script):
"""
:param cache_file: Путь к файлу кеша (например, "fail2ban_cache.json").
:param monitoring_script: Путь к shell‑скрипту (например, "./fail2ban_parse.sh").
Скрипт должен выводить валидный JSON‑массив объектов с ключами:
- "service" (например, "guacamole")
- "time" (формат "YYYY-MM-DD HH:MM:SS")
- "ip"
- "action" (в данном случае всегда "Ban")
"""
self.cache_file = cache_file
self.monitoring_script = monitoring_script
def load_cache(self):
"""
Загружает кеш из файла.
Структура кеша: {
"ip1": [event1, event2, ...],
"ip2": [...],
...
}
При загрузке остаются только события текущего дня.
"""
if not os.path.exists(self.cache_file):
return {}
try:
with open(self.cache_file, "r") as f:
cache = json.load(f)
current_day = datetime.now().strftime("%Y-%m-%d")
new_cache = {}
for ip, events in cache.items():
filtered = [event for event in events if event.get("time") and parse_event_day(event["time"]) == current_day]
if filtered:
new_cache[ip] = filtered
return new_cache
except Exception as e:
logger.error(f"Ошибка при загрузке кеша: {e}")
return {}
def save_cache(self, cache):
"""Сохраняет кеш в файл."""
try:
with open(self.cache_file, "w") as f:
json.dump(cache, f, indent=4)
except Exception as e:
logger.error(f"Ошибка при сохранении кеша: {e}")
def process_script_output(self):
"""
Выполняет shell‑скрипт для парсинга логов Fail2Ban.
Ожидается, что скрипт выводит валидный JSON‑массив объектов.
"""
try:
result = subprocess.run(
[self.monitoring_script],
capture_output=True,
text=True,
check=True
)
output = result.stdout.strip()
data = json.loads(output)
return data
except json.JSONDecodeError as e:
logger.error(f"Скрипт вернул невалидный JSON. Ошибка: {e}")
return []
except Exception as e:
logger.error(f"Ошибка выполнения скрипта: {e}")
return []
def update_cache(self, new_events, cache):
"""
Обновляет кеш с данными о событиях.
Ключ кеша – IP, а для каждого IP хранится массив событий.
Для каждого нового события (уникальность определяется по комбинации time и action)
если такого события ещё нет в кеше, оно добавляется и считается новым.
"""
diff = [] # список новых событий
current_day = datetime.now().strftime("%Y-%m-%d")
for event in new_events:
if "time" not in event:
continue
if parse_event_day(event['time']) != current_day:
continue
ip = event['ip']
if ip not in cache:
cache[ip] = []
unique_key = f"{event['time']}|{event['action']}"
exists = any(f"{e.get('time')}|{e.get('action')}" == unique_key for e in cache[ip])
if not exists:
cache[ip].append(event)
diff.append(event)
return diff, cache
def produce_notification(self, diff):
"""
Генерирует уведомления для новых событий Fail2Ban в Markdown‑формате.
Для каждого нового события выводятся:
- Service
- Time
- IP
- Action
"""
if not diff:
return ""
message = ""
for event in diff:
message += f"📢 Fail2Ban Alert\n"
message += "```\n"
message += f"Service: {event.get('service', 'unknown')}\n"
message += f"Time : {event.get('time', 'N/A')}\n"
message += f"IP : {event.get('ip', 'N/A')}\n"
message += f"Action : {event.get('action', 'N/A')}\n"
message += "```\n\n"
return message
def run_monitoring(self):
"""
Основной процесс мониторинга:
1. Загружает кеш.
2. Получает новые события из скрипта.
3. Обновляет кеш, определяя новые события.
4. Сохраняет обновлённый кеш.
5. Генерирует уведомление, если обнаружены новые события.
"""
cache = self.load_cache()
new_events = self.process_script_output()
if not new_events:
#logger.error("Не удалось получить данные из скрипта или данные отсутствуют.")
return 0
diff, cache = self.update_cache(new_events, cache)
self.save_cache(cache)
notification_message = self.produce_notification(diff)
if notification_message:
logger.info("Обнаружены новые события Fail2Ban:\n" + notification_message)
return notification_message
else:
logger.debug("Новых событий Fail2Ban не обнаружено.")
return 0
module.py 11.3 КБ
"""network.Module — сканер сети + wireguard/guacamole/f2b-мониторы.
Сканер (NetworkScanner) собирается здесь и стартует фоном в start()
(daemon-поток, как старый messengers/telegram.py). Отправка — через
sender активного транспорта (берётся в start, транспорты уже подняты).
Права — RightsCheckerAdapter(ctx.rights). Кеши — data/cache/network
(пишет сам; прод-кеш не трогаем). macblock TG-кнопки — legacy через ui_tg.
"""
import logging
import os
import threading
from core.base_module import BaseModule
from core.legacy import RightsCheckerAdapter
from . import config as pkg_config
logger = logging.getLogger(__name__)
class Module(BaseModule):
def setup(self, ctx) -> None:
mod_dir = os.path.dirname(os.path.abspath(__file__))
self.config = ctx.module_config("network", mod_dir)
self.ctx = ctx
net_dir = os.path.join(ctx.data_dir, "cache", "network")
def _f(key, default):
return (self.config.get(key, "") or "").strip() or default
def _fp(key, default):
from core.paths import abs_data
return abs_data(ctx.data_dir,
(self.config.get(key, "") or "").strip(),
default)
files = {
"MESSAGE_QUEUE_FILE": _fp(
"MESSAGE_QUEUE_FILE", os.path.join(net_dir, "message_queue.json")),
"CACHE_FILE_NETWORK": _fp(
"CACHE_FILE_NETWORK", os.path.join(net_dir, "networkscan_cache.json")),
"CACHE_FILE_WG": _fp(
"CACHE_FILE_WG", os.path.join(net_dir, "networkscan_cache_wg.json")),
# wg/guacamole-скрипты живут в своих модулях (файлы на месте
# независимо от загрузки модулей).
"SCRIPT_WG_DUMP": os.path.join(
os.path.dirname(mod_dir), "network_wireguard", "scripts", "wg_dump.sh"),
"SCRIPT_WG_PING": os.path.join(
os.path.dirname(mod_dir), "network_wireguard", "scripts", "wg_ping.sh"),
"SCRIPT_AG_LOGIN": os.path.join(
os.path.dirname(mod_dir), "network_guacamole", "scripts", "ag_login.sh"),
"SCRIPT_AG_CACHE": _fp(
"SCRIPT_AG_CACHE", os.path.join(net_dir, "guac_cache.json")),
"SCRIPT_F2B_CHECK": os.path.join(mod_dir, "scripts_net", "f2b_check.sh"),
"SCRIPT_F2B_CACHE": _fp(
"SCRIPT_F2B_CACHE", os.path.join(net_dir, "f2b_cache.json")),
"SCAN_NETWORK_IGNORED_MACS": _fp(
"SCAN_NETWORK_IGNORED_MACS", os.path.join(net_dir, "networkscan_ignored.json")),
}
# Каталоги кешей обязан создать модуль (сканер пишет без makedirs).
for _d in (net_dir,):
try:
os.makedirs(_d, exist_ok=True)
except Exception: # noqa: BLE001
pass
# install ДО импорта движков: macblock биндит значения при импорте.
pkg_config.install(dict(self.config.values), files)
try:
from . import scan as scan_mod
except ImportError as e:
from core.errors import ConfigError
raise ConfigError("network: движки требуют venv (telegram): %s" % e)
self.scan_mod = scan_mod
self.rights_adapter = RightsCheckerAdapter(ctx.rights)
f2b_on = (self.config.get("F2B_ENABLED", "1") or "1").strip() not in ("0", "no", "false")
# Цели уведомлений — единый ADMIN_IDS (модуль access, без дублей).
try:
notify_ids = ctx.api.call("access.admin_ids")
except Exception: # noqa: BLE001
notify_ids = []
self.scanner = scan_mod.NetworkScanner(
files["MESSAGE_QUEUE_FILE"], files["CACHE_FILE_NETWORK"],
files["CACHE_FILE_WG"], files["SCRIPT_WG_DUMP"], files["SCRIPT_WG_PING"],
files["SCRIPT_AG_LOGIN"], files["SCRIPT_AG_CACHE"], f2b_on,
files["SCRIPT_F2B_CHECK"], files["SCRIPT_F2B_CACHE"],
notify_ids, self.rights_adapter, None,
pkg_config.SCAN_NETWORK_RANGE, pkg_config.CHECK_DURATION_NSCAN_LONG,
pkg_config.CHECK_DURATION_NSCAN_FAST, files["SCAN_NETWORK_IGNORED_MACS"],
on_diff=self._on_diff)
self._thread = None
self._mon_thread = None
self._mon_stop = threading.Event()
self._monitor_registry = {}
self._files = files
self._f2b_on = f2b_on
ctx.api.register("network", "scanner", lambda: self.scanner)
ctx.api.register("network", "cache_file", lambda: files["CACHE_FILE_NETWORK"])
ctx.api.register(
"network", "ignored_file", lambda: files["SCAN_NETWORK_IGNORED_MACS"])
ctx.api.register("network", "register_monitor", self._register_monitor)
from . import store as _store
ctx.api.register("network", "devices_snapshot", _store.devices_snapshot)
def _on_diff(self, diff):
"""Хук сканера: diff -> события шины (по одному на устройство)."""
mapping = (
("new", "network.device_new"),
("reconnected", "network.device_reconnected"),
("disconnected", "network.device_disconnected"),
("ip_changed", "network.device_ip_changed"),
)
for section, event in mapping:
for mac, dev in ((diff or {}).get(section) or {}).items():
try:
self.ctx.events.emit(event, {
"mac": mac,
"ip": dev.get("ip"),
"old_ip": dev.get("old_ip"),
"name": dev.get("name"),
"name_custom": dev.get("name_custom") or "",
"type": dev.get("type"),
"ports": list(dev.get("ports") or []),
"first_appearance": dev.get("first_appearance"),
"last_appearance": dev.get("last_appearance"),
"last_disconnection": dev.get("last_disconnection"),
"connections_total": dev.get("connections_total"),
})
except Exception: # noqa: BLE001
logger.exception("network emit %s failed", event)
def _register_monitor(self, name, spec):
"""Плагин доп. источника: {factory, interval_s, type, silent}.
factory() -> объект с run_monitoring() -> str. Вызывают network_*
модули в своём setup; интервалы у каждого свои (свой .env).
"""
if not name or not isinstance(spec, dict):
raise ValueError("register_monitor: нужно имя и spec")
factory = spec.get("factory")
if not callable(factory):
raise ValueError("register_monitor %s: нет factory" % name)
try:
interval = int(spec.get("interval_s", 60) or 60)
except (TypeError, ValueError):
raise ValueError("register_monitor %s: interval_s — число" % name)
interval = max(5, min(86400, interval))
notif_type = str(spec.get("type") or name)
self._monitor_registry[str(name)] = {
"factory": factory, "interval": interval,
"type": notif_type, "silent": bool(spec.get("silent", False)),
}
logger.info("network: монитор %s каждые %s с", name, interval)
return True
def start(self) -> None:
try:
self.scanner.bot = self.ctx.api.call("transport_telegram.sender")
except Exception as e: # noqa: BLE001
logger.warning("network: нет sender: %s", e)
# Штатный f2b-монитор сканера + зарегистрированные network_*.
from . import fail2ban_mon as _f2b
monitors = {"f2b": {
"monitor": _f2b.Fail2BanMonitor(
self._files["SCRIPT_F2B_CACHE"],
self._files["SCRIPT_F2B_CHECK"],
) if self._f2b_on else None,
"interval": int(self.config.get("CHECK_DURATION_NSCAN_FAST", "5") or 5),
"type": "f2b", "silent": True}}
for name, spec in self._monitor_registry.items():
try:
monitors[name] = {
"monitor": spec["factory"](),
"interval": spec["interval"],
"type": spec["type"], "silent": spec["silent"]}
except Exception as e: # noqa: BLE001
logger.warning("network: монитор %s не поднялся: %s", name, e)
self._monitors = {n: m for n, m in monitors.items()
if m.get("monitor") is not None}
self._mon_stop.clear()
self._thread = threading.Thread(target=self._run_scan, daemon=True,
name="net-scanner")
self._thread.start()
self._mon_thread = threading.Thread(target=self._run_monitors, daemon=True,
name="net-monitors")
self._mon_thread.start()
def _run_monitors(self):
import time as _time
due = {}
tick = 5
try:
tick = max(1, min([m["interval"] for m in self._monitors.values()] or [5]))
except Exception: # noqa: BLE001
pass
while not self._mon_stop.wait(tick):
now = _time.time()
for name, spec in list(self._monitors.items()):
if now < due.get(name, 0):
continue
due[name] = now + spec["interval"]
try:
msg = spec["monitor"].run_monitoring()
except Exception as e: # noqa: BLE001
logger.warning("network monitor %s: %s", name, type(e).__name__)
continue
try:
if msg:
self.scanner.send_notification(
msg, notif_type=spec["type"], disable_notif=spec["silent"])
else:
with self.scanner._queue_lock:
queued = [m for m in self.scanner.load_message_queue()
if m.get("type") == spec["type"]]
if queued:
self.scanner.send_notification(
"", notif_type=spec["type"], disable_notif=spec["silent"])
except Exception as e: # noqa: BLE001
logger.warning("network notify %s: %s", name, e)
def stop(self) -> None:
self._mon_stop.set()
for t in (self._thread, self._mon_thread):
if t is not None:
try:
t.join(timeout=10)
except Exception: # noqa: BLE001
pass
def _run_scan(self):
try:
self.scanner.start_scan()
except Exception as e: # noqa: BLE001
logger.error("network scan: %s", e)
def health(self) -> dict:
alive = self._thread.is_alive() if self._thread else False
return {"ok": self.scanner is not None, "module": self.name,
"scanner": alive}
scan.py 38.4 КБ
import subprocess
import json
import time
import logging
import os
from datetime import datetime, timedelta
import threading
from threading import Thread
from telegram import InlineKeyboardButton, InlineKeyboardMarkup
from .fail2ban_mon import Fail2BanMonitor # BOTMOD: штатный монитор сканера
import re
import copy
from . import config # BOTMOD: было import config
logger = logging.getLogger(__name__)
def escape_markdown_v2(text):
"""Экранирует специальные символы для MarkdownV2.
Содержимое fenced-блоков кода (между ```) не экранируется,
иначе символы внутри блоков отображались бы с лишними обратными слэшами.
"""
in_code = False
result = []
for line in text.split("\n"):
if line.strip().startswith("```"):
in_code = not in_code
result.append(line)
continue
if in_code:
result.append(line)
else:
result.append(re.sub(r'([_*\[\]()~>#+\-=|{}.!])', r'\\\1', line))
return "\n".join(result)
def format_duration(seconds):
return str(timedelta(seconds=seconds))
def is_host_reachable(ip):
"""
Проверяет, доступен ли хост через пинг.
Возвращает True, если хотя бы один пакет успешно дошёл.
"""
try:
# Подробный пинг (10 пакетов) важен для проверки "отключаемых" устройств:
# мобильные устройства (телефоны) могут не всегда отвечать на nmap-скан,
# поэтому чтобы не считать ложное отключение, даём им достаточно попыток.
result = subprocess.run(
["ping", "-c", "10", ip],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
timeout=20
)
# Успех, если код возврата 0
return result.returncode == 0
except Exception as e:
logger.error(f"Ошибка проверки пинга для {ip}: {e}")
return False
def _extract_new_device_macs(text):
"""Извлекает MAC-адреса новых устройств (секции, начинающиеся с 🚨)."""
out = []
for m in re.finditer(r"\U0001F6A8.{0,300}?MAC:\s*((?:[0-9A-Fa-f]{2}[:-]){5}[0-9A-Fa-f]{2})",
text, re.S):
mac = m.group(1).upper().replace("-", ":")
if mac not in out:
out.append(mac)
return out
# Лимит Telegram — 4096 символов, но MarkdownV2-экранирование раздувает
# текст (каждый спецсимвол + обратный слэш). Режем с запасом: 3000 сырых.
TG_CHUNK_LIMIT = 3000
def split_tg_message(text, limit=TG_CHUNK_LIMIT):
"""Режет длинное уведомление на части по границам блоков (пустая строка).
Блоки устройств (~150 символов) никогда не рвутся пополам; одиночный
сверхдлинный блок без '\n\n' режется жёстко по лимиту.
"""
text = text or ""
if len(text) <= limit:
return [text]
parts, cur, cur_len = [], [], 0
for block in text.split("\n\n"):
b = block.strip("\n")
if not b:
continue
if len(b) > limit:
# аварийный случай: сам блок больше лимита — режем жёстко
if cur:
parts.append("\n\n".join(cur))
cur, cur_len = [], 0
for i in range(0, len(b), limit):
parts.append(b[i:i + limit])
continue
if cur and cur_len + len(b) + 2 > limit:
parts.append("\n\n".join(cur))
cur, cur_len = [], 0
cur.append(b)
cur_len += len(b) + 2
if cur:
parts.append("\n\n".join(cur))
return parts or [text]
class NetworkScanner:
def __init__(self, message_queue, cache_file, wg_cache_file, monitoring_script, ping_script,
ag_script, ag_cache, f2b_enabled, f2b_script, f2b_cache, admin_ids, rights_checker_class, sender,
network_range, long_check_interval, fast_check_interval, mac_filter_file,
on_diff=None):
self.message_queue_file = message_queue # путь к файлу очереди сообщений
self._queue_lock = threading.Lock() # защищает очередь сообщений от гонок потоков
self.cache_file = cache_file
self.admin_ids = admin_ids
self.rights_checker = rights_checker_class
# sender — шлюз messengers.sender.GatewaySender (активный транспорт).
self.bot = sender
self.network_range = network_range
self.long_check_interval = long_check_interval # интервал для долгих проверок (сканирование сети)
self.fast_check_interval = fast_check_interval # интервал для быстрых проверок (мониторинг доп. источников)
self.mac_filter_file = mac_filter_file
self.excluded_macs = self.load_mac_filter()
self.wg_monitor = None # BOTMOD: мониторы доп. источников вешает модуль (register_monitor)
self.ag_monitor = None
self.f2b_monitor = Fail2BanMonitor(f2b_cache, f2b_script)
# on_diff(diff) — хук сырых diff скана для шины событий (ставит модуль).
self._on_diff = on_diff
def load_cache(self):
"""Загружает кэш из JSON-файла"""
if not os.path.exists(self.cache_file):
return {}
with open(self.cache_file, 'r') as f:
try:
return json.load(f)
except Exception as e:
logger.error(f"Ошибка загрузки кэша: {e}")
return {}
def save_cache(self, cache):
"""Сохраняет кэш в JSON-файл"""
with open(self.cache_file, 'w') as f:
json.dump(cache, f, indent=4)
def load_mac_filter(self):
"""Загружает список MAC-адресов, которые нужно исключить, из JSON файла."""
if not os.path.exists(self.mac_filter_file):
return []
try:
with open(self.mac_filter_file, 'r') as f:
mac_list = json.load(f)
if not isinstance(mac_list, list):
logger.error("Формат JSON для фильтра MAC должен быть списком.")
return []
# Приводим все MAC-адреса к верхнему регистру для корректного сравнения
return [mac.upper() for mac in mac_list]
except Exception as e:
logger.error(f"Ошибка загрузки MAC фильтра: {e}")
return []
def perform_network_scan(self):
"""
Выполняет sweep-скан сети с помощью nmap:
nmap -sn -T4 --privileged <NETWORK_RANGE>
Возвращает словарь устройств.
"""
cmd = ['nmap', '-sn', '-T4', '--privileged', self.network_range]
try:
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
except subprocess.CalledProcessError as e:
logger.error(f"Ошибка сканирования сети: {e}")
return {}
devices = {}
current_device = None
self.excluded_macs = self.load_mac_filter()
for line in result.stdout.splitlines():
line = line.strip()
if line.startswith("Nmap scan report for"):
if '(' in line and ')' in line:
ip = line.split('(')[-1].split(')')[0]
hostname = line.split('(')[0].split()[-1]
else:
ip = line.split()[-1]
hostname = "Unknown"
current_device = {
"ip": ip,
"name": hostname,
"mac": None,
"ports": [],
"type": "Unknown"
}
elif line.startswith("MAC Address:") and current_device is not None:
parts = line.split()
mac = parts[2].upper()
if mac in self.excluded_macs:
current_device = None
continue
current_device["mac"] = mac
if "(" in line and ")" in line:
vendor = line.split("(")[-1].split(")")[0]
current_device["type"] = vendor.strip()
devices[current_device["mac"]] = current_device
current_device = None
return devices
def perform_port_scan(self, ip):
"""
Выполняет полный портскан для указанного IP:
nmap -p- --open -T4 --privileged <ip>
Возвращает список открытых портов.
"""
cmd = ['nmap', '-p-', '--open', '-T5', ip]
try:
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
except subprocess.CalledProcessError as e:
logger.error(f"Ошибка сканирования портов для {ip}: {e}")
return []
ports = []
for line in result.stdout.splitlines():
if "/tcp" in line and "open" in line:
port = line.split("/")[0]
ports.append(port)
return ports if ports else ["Нет открытых портов"]
def add_ignored_mac(self, mac):
mac = mac.upper()
cache = self.load_cache()
if mac not in cache:
return False, "MAC адрес не найден в кэше сети."
ignored = self.load_mac_filter()
if mac in ignored:
return False, "MAC адрес уже присутствует в игнор-листе."
ignored.append(mac)
try:
with open(self.mac_filter_file, "w", encoding="utf-8") as f:
json.dump(ignored, f, indent=4)
return True, "MAC адрес успешно добавлен в игнор-лист."
except Exception as e:
return False, f"Ошибка сохранения: {e}"
def remove_ignored_mac(self, mac):
mac = mac.upper()
ignored = self.load_mac_filter()
if mac not in ignored:
return False, "MAC адрес отсутствует в игнор-листе."
ignored.remove(mac)
try:
with open(self.mac_filter_file, "w", encoding="utf-8") as f:
json.dump(ignored, f, indent=4)
return True, "MAC адрес успешно удалён из игнор-листа."
except Exception as e:
return False, f"Ошибка сохранения: {e}"
def _ensure_stats(self, dev, now):
"""
Гарантирует наличие полей статистики у устройства и выполняет сброс
суточных счётчиков при смене календарного дня.
"""
if "connections_total" not in dev:
dev["connections_total"] = 0
if "connections_today" not in dev:
dev["connections_today"] = 0
if "online_seconds_total" not in dev:
dev["online_seconds_total"] = 0
if "online_seconds_today" not in dev:
dev["online_seconds_today"] = 0
if "online_start_ts" not in dev:
dev["online_start_ts"] = None
today = datetime.now().strftime("%Y-%m-%d")
if dev.get("stats_date") != today:
dev["stats_date"] = today
dev["connections_today"] = 0
dev["online_seconds_today"] = 0
if dev.get("connected") and dev.get("online_start_ts") is None:
dev["online_start_ts"] = now
return dev
def update_cache_and_diff(self, scanned_devices):
"""
Сравнивает данные текущего сканирования с кэшем, обновляет статус `connected`
и возвращает diff, при этом пытаясь сохранить изменения (например, индивидуальное
кастомное имя) внесённые параллельно другим процессом.
"""
# Шаг 1. Загружаем исходный кэш и делаем глубокую копию для редактирования.
baseline_cache = self.load_cache() # эталонная версия
editable_cache = copy.deepcopy(baseline_cache) # версия для редактирования
now = time.time()
diff = {"new": {}, "reconnected": {}, "disconnected": {}, "ip_changed": {}}
# Шаг 2. Обработка данных сканирования (scanned_devices)
for mac, dev in scanned_devices.items():
if mac not in editable_cache:
dev["first_appearance"] = now
dev["last_appearance"] = now
dev["last_disconnection"] = None
dev["ports"] = self.perform_port_scan(dev["ip"])
dev["connected"] = True
dev["name_custom"] = ""
dev["miss_count"] = 0
dev["stats_date"] = datetime.now().strftime("%Y-%m-%d")
dev["connections_total"] = 1
dev["connections_today"] = 1
dev["online_seconds_total"] = 0
dev["online_seconds_today"] = 0
dev["online_start_ts"] = now
editable_cache[mac] = dev
diff["new"][mac] = dev
else:
cached = editable_cache[mac]
# Сохраняем ранее установленные значения, в том числе кастомное имя.
dev["first_appearance"] = cached.get("first_appearance", now)
dev["last_disconnection"] = cached.get("last_disconnection")
dev["name_custom"] = cached.get("name_custom", "")
# Устройство видно в сканировании — сбрасываем счётчик пропусков.
dev["miss_count"] = 0
# Переносим накопленную статистику из кэша и делаем суточный сброс при необходимости.
dev["connections_total"] = cached.get("connections_total", 0)
dev["connections_today"] = cached.get("connections_today", 0)
dev["online_seconds_total"] = cached.get("online_seconds_total", 0)
dev["online_seconds_today"] = cached.get("online_seconds_today", 0)
dev["online_start_ts"] = cached.get("online_start_ts")
dev["stats_date"] = cached.get("stats_date")
self._ensure_stats(dev, now)
# Определяем смену IP-адреса (например, при обновлении DHCP-аренды).
old_ip = cached.get("ip")
if old_ip and old_ip != dev.get("ip"):
ip_info = copy.deepcopy(dev)
ip_info["old_ip"] = old_ip
diff["ip_changed"][mac] = ip_info
if not cached.get("connected"):
dev["last_appearance"] = now
dev["ports"] = self.perform_port_scan(dev["ip"])
dev["connected"] = True
dev["connections_total"] += 1
dev["connections_today"] += 1
dev["online_start_ts"] = now
editable_cache[mac] = dev
diff["reconnected"][mac] = dev
else:
dev["last_appearance"] = cached.get("last_appearance", now)
dev["ports"] = cached.get("ports", [])
dev["connected"] = True
editable_cache[mac] = dev
# Обработка устройств, которых нет в текущем сканировании.
# Отключённым считаем только после DEBOUNCE_COUNT пропусков подряд
# (нет в nmap И не отвечает на пинг), чтобы не срабатывать на ложные пропадания.
for mac in list(editable_cache.keys()):
if mac not in scanned_devices and editable_cache[mac].get("connected"):
dev = editable_cache[mac]
self._ensure_stats(dev, now)
ip = dev.get("ip")
reachable = bool(ip) and is_host_reachable(ip)
if reachable:
dev["miss_count"] = 0
else:
miss = dev.get("miss_count", 0) + 1
dev["miss_count"] = miss
if miss >= config.DEBOUNCE_COUNT:
# Начисляем время онлайн и закрываем интервал.
online_start = dev.get("online_start_ts")
if online_start:
dev["online_seconds_total"] = dev.get("online_seconds_total", 0) + (now - online_start)
today_start = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0).timestamp()
dev["online_seconds_today"] = dev.get("online_seconds_today", 0) + (now - max(online_start, today_start))
dev["online_start_ts"] = None
dev["last_disconnection"] = now
dev["connected"] = False
diff["disconnected"][mac] = dev
# Шаг 3. Прежде чем сохранить, пытаемся сохранить изменения, внесённые параллельно.
current_cache = self.load_cache() # Текущая версия кэша из файла
if current_cache != baseline_cache:
# Для каждого MAC-адреса проверяем, если поле name_custom изменилось в current_cache,
# то обновляем его в editable_cache.
for mac, current_dev in current_cache.items():
if mac in baseline_cache:
# Если значение name_custom в текущем кэше отличается от эталонного,
# то используем его в editable_cache.
baseline_name = baseline_cache[mac].get("name_custom", "")
current_name = current_dev.get("name_custom", "")
if current_name != baseline_name:
# Обновляем кастомное имя в редактируемом кэше
if mac in editable_cache:
editable_cache[mac]["name_custom"] = current_name
else:
# Если новое устройство появилось в current_cache, добавляем его в editable_cache.
self._ensure_stats(current_dev, now)
editable_cache[mac] = current_dev
# Шаг 4. Сохраняем полностью обновлённый кэш и возвращаем сформированный diff.
self.save_cache(editable_cache)
return diff
def produce_notification(self, diff):
"""Формирует сообщение об изменениях в сети."""
message = ""
# Новые устройства
if diff.get("new"):
for mac, dev in diff["new"].items():
rule = "new"
ip = dev.get("ip", "Unknown")
message += f"🚨 {ip} ({rule}):\n"
message += "```\n"
message += f"MAC: {mac}\n"
message += f"Name: {dev.get('name', 'Unknown')}\n"
message += f"Ports: {', '.join(dev.get('ports', []))}\n"
message += f"Type: {dev.get('type', 'Unknown')}\n"
message += f"First appearance: {datetime.fromtimestamp(dev.get('first_appearance')).strftime('%Y-%m-%d %H:%M:%S')}\n"
message += "```\n\n"
# Повторное подключение
if diff.get("reconnected"):
for mac, dev in diff["reconnected"].items():
last_disconnection = dev.get("last_disconnection")
name_custom = dev.get("name_custom", "")
ip = dev.get("ip", "Unknown")
if last_disconnection is not None:
duration = format_duration(int(time.time() - last_disconnection))
rule = f"connected after {duration} offline"
else:
rule = "connected"
if name_custom:
message += f"🔄 {name_custom} ({rule}):\n"
else:
message += f"🔄 {ip} ({rule}):\n"
message += "```\n"
if name_custom:
message += f"IP: {ip}\n"
message += f"MAC: {mac}\n"
message += f"Name: {dev.get('name', 'Unknown')}\n"
message += f"Ports: {', '.join(dev.get('ports', []))}\n"
message += f"Type: {dev.get('type', 'Unknown')}\n"
message += f"First appearance: {datetime.fromtimestamp(dev.get('first_appearance')).strftime('%Y-%m-%d %H:%M:%S')}\n"
message += "```\n\n"
# Отключившиеся устройства
if diff.get("disconnected"):
for mac, dev in diff["disconnected"].items():
last_appearance = dev.get("last_appearance")
name_custom = dev.get("name_custom", "")
ip = dev.get("ip", "Unknown")
if last_appearance is not None:
duration = format_duration(int(time.time() - last_appearance))
rule = f"disconnected after {duration} online"
else:
rule = "disconnected"
if name_custom:
message += f"🔌 {name_custom} ({rule}):\n"
else:
message += f"🔌 {ip} ({rule}):\n"
message += "```\n"
if name_custom:
message += f"IP: {ip}\n"
message += f"MAC: {mac}\n"
message += f"Name: {dev.get('name', 'Unknown')}\n"
message += f"Ports: {', '.join(dev.get('ports', []))}\n"
message += f"Type: {dev.get('type', 'Unknown')}\n"
message += f"First appearance: {datetime.fromtimestamp(dev.get('first_appearance')).strftime('%Y-%m-%d %H:%M:%S')}\n"
message += "```\n\n"
# Смена IP-адреса
if diff.get("ip_changed"):
for mac, dev in diff["ip_changed"].items():
name_custom = dev.get("name_custom", "")
ip = dev.get("ip", "Unknown")
old_ip = dev.get("old_ip", "Unknown")
title = name_custom if name_custom else ip
message += f"📍 {title} сменил IP:\n"
message += f" • Было: {old_ip}\n"
message += f" • Стало: {ip}\n"
message += f" • MAC: {mac}\n"
message += f" • Имя: {dev.get('name', 'Unknown')}\n\n"
return message
def load_message_queue(self):
"""Загружает очередь сообщений из файла.
Каждая запись – словарь с ключами: "timestamp", "message" и "type".
"""
if not os.path.exists(self.message_queue_file):
return []
try:
with open(self.message_queue_file, "r") as f:
return json.load(f)
except Exception as e:
logger.error(f"Ошибка загрузки очереди сообщений: {e}")
return []
def save_message_queue(self, queue):
"""Сохраняет очередь сообщений в JSON-файл."""
try:
with open(self.message_queue_file, "w") as f:
json.dump(queue, f, indent=4)
except Exception as e:
logger.error(f"Ошибка сохранения очереди сообщений: {e}")
def generate_inline_buttons_from_raw_message(self, raw_message):
"""
Парсит MAC-адреса из сырого сообщения, загружает кэш из файла и сравнивает с ним.
Для каждого MAC, присутствующего в кэше и у которого поле name_custom пустое,
формирует кнопку для добавления/изменения имени.
Возвращает объект InlineKeyboardMarkup или None, если кнопок нет.
"""
try:
# Используем регулярное выражение для поиска MAC-адресов (форматы XX:XX:XX:XX:XX:XX или XX-XX-XX-XX-XX-XX).
mac_pattern = r'(?:[0-9A-Fa-f]{2}[:\-]){5}[0-9A-Fa-f]{2}'
mac_addresses = re.findall(mac_pattern, raw_message)
if not mac_addresses:
return None
network_cache = self.load_cache()
add_buttons = []
for mac in mac_addresses:
mac_upper = mac.upper()
if mac_upper in network_cache:
dev = network_cache[mac_upper]
# Если поле name_custom пустое, добавляем кнопку
if not dev.get("name_custom", "").strip():
add_buttons.append(InlineKeyboardButton(f"📝 {mac_upper}", callback_data=f"nmap_setname:{mac_upper}"))
if add_buttons:
# Группируем кнопки по 2 в ряд (можно настроить количество столбцов)
rows = [add_buttons[i:i+2] for i in range(0, len(add_buttons), 2)]
return InlineKeyboardMarkup(rows)
else:
return None
except Exception as e:
logger.error(f"Ошибка формирования inline-кнопок по сырому сообщению: {e}")
return None
def _send_with_retry(self, chat_id, text, **kwargs):
"""
Отправляет сообщение с ретраями (NOTIF_MAX_RETRIES) и экспоненциальной
задержкой между попытками (NOTIF_RETRY_BACKOFF). Возвращает True при успехе.
"""
delay = config.NOTIF_RETRY_BACKOFF
for attempt in range(1, config.NOTIF_MAX_RETRIES + 1):
try:
self.bot.send_message(chat_id=chat_id, text=text, **kwargs)
return True
except Exception as e:
logger.warning(
f"Отправка сообщения для {chat_id}: попытка {attempt}/{config.NOTIF_MAX_RETRIES} не удалась: {e}"
)
if attempt < config.NOTIF_MAX_RETRIES:
time.sleep(delay)
delay *= 2
return False
def send_notification(self, raw_message, notif_type, disable_notif=False):
"""
Отправляет уведомление определённого типа (notif_type) администраторам.
Все уведомления одного типа отправляются отдельно – никакого комбинирования разных типов.
Для каждого администратора проверяется наличие права:
- для типа "wireguard" – receive_notifications_wireguard,
- для "guacamole" – receive_notifications_guacamole,
- для "f2b" – receive_notifications_f2b,
- для "status" – receive_notifications_status.
Параметр disable_notif (или настройка receive_notifications_no_sound) влияет на флаг disable_notification.
"""
# Определяем требуемое право для данного типа уведомления
permission_mapping = {
"wireguard": "receive_notifications_wireguard",
"guacamole": "receive_notifications_guacamole",
"f2b": "receive_notifications_f2b",
"status": "receive_notifications_status",
"nmap": "receive_notifications_nmap"
}
permission_key = permission_mapping.get(notif_type)
if not permission_key:
logger.error(f"Неизвестный тип уведомлений: {notif_type}")
return
# Весь блок под локом: очереди читаются/пишутся из двух потоков (быстрые и долгие
# проверки), поэтому чтение->отправка->сохранение должны быть атомарными,
# иначе возможна потеря или дублирование уведомлений.
with self._queue_lock:
# Загружаем очередь и фильтруем записи по notif_type
queue = self.load_message_queue()
filtered_queue = [msg for msg in queue if msg.get("type") == notif_type]
# Формируем список сообщений для отправки. Длинные (агрегаты nmap,
# очередь) режем на части по 3000 символов — иначе Telegram
# отвечает 'Message is too long', ретраи спамят в лог, а текст
# вечно перекладывается в очередь и не доходит.
messages = []
if raw_message:
messages.append(raw_message)
for queued_msg in filtered_queue:
messages.append(queued_msg["message"] + " [Отправлено с задержкой, первоначальное время: " + queued_msg["timestamp"] + "]\n\n")
if not messages:
return # нечего отправлять
# Разбиение: [(чанк, последний_ли)] — клавиатуры вешаем только
# на последний чанк каждого логического сообщения.
outgoing = []
for message in messages:
chunks = split_tg_message(message)
for ci, ch in enumerate(chunks):
outgoing.append((ch, ci == len(chunks) - 1))
# Новая логика для типа nmap: добавляем inline-клавиатуру с кнопками "Добавить имя"
additional_markup = self.generate_inline_buttons_from_raw_message("\n".join(messages))
# Кнопки блокировки для новых устройств (право network_manage)
acl_block_rows = []
if notif_type == "nmap":
new_macs = _extract_new_device_macs("\n".join(messages))
for i in range(0, len(new_macs), 2):
acl_block_rows.append([
InlineKeyboardButton(
f"📔 Блокировать {m}",
callback_data=f"acl_ask:{m}")
for m in new_macs[i:i+2]
])
logger.info(f"Отправка уведомления типа {notif_type}: "
f"{len(messages)} сообщ., {len(outgoing)} чанков, "
f"всего {sum(len(m) for m in messages)} символов")
any_failed = False
for admin_id in self.admin_ids:
can_receive, no_sound = self.rights_checker.should_notify(admin_id, permission_key)
if not can_receive:
continue
try:
# Права — только через checker (ns-aware: ищет
# 'telegram:<id>'). Прямой rights.get(str_id) мимо
# ключей 'telegram:*' — так клавиатура не появлялась
# вообще никогда.
has = self.rights_checker.has_right
# Определяем disable_notification: совмещаем переданный параметр и настройку "без звука"
send_disable = disable_notif or no_sound
if has(admin_id, "nmap_manage") or has(admin_id, "all"):
reply_markup_sent = additional_markup
else:
reply_markup_sent = None
# Кнопки блокировки новых устройств — только с правом network_manage
if acl_block_rows and (has(admin_id, "network_manage") or has(admin_id, "all")):
base_rows = list(reply_markup_sent.inline_keyboard) if reply_markup_sent else []
reply_markup_sent = InlineKeyboardMarkup(base_rows + acl_block_rows)
for chunk, is_last in outgoing:
ok = self._send_with_retry(
admin_id,
escape_markdown_v2(chunk),
disable_notification=send_disable,
parse_mode='MarkdownV2',
reply_markup=reply_markup_sent if is_last else None
)
if not ok:
any_failed = True
except Exception as e:
logger.warning(f"Ошибка отправки уведомления для {admin_id}: {e}")
any_failed = True
if any_failed and raw_message:
# Если не удалось доставить – кладём сообщение в очередь, избегая дублей.
unsent = self.load_message_queue()
if not any(m.get("type") == notif_type and m.get("message") == raw_message for m in unsent):
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
unsent.append({"timestamp": timestamp, "message": raw_message, "type": notif_type})
self.save_message_queue(unsent)
elif not any_failed:
# Все доставлено успешно – очищаем из очереди записи, относящиеся к данному типу
new_queue = [msg for msg in queue if msg.get("type") != notif_type]
self.save_message_queue(new_queue)
def start_long_checks(self):
"""
Долгие проверки (сканирование сети с помощью nmap)
выполняются с интервалом self.long_check_interval.
Важное событие (появление нового устройства) отправляется сразу.
Неважные события (переподключения/отключения) копятся в буфере и отправляются
одним сообщением после накопления AGGREGATION_MAX_CHARS символов или по
истечении AGGREGATION_INTERVAL секунд с первого события в буфере.
"""
pending = []
pending_first_ts = None
while True:
logger.debug("Запуск долгих проверок: сканирование сети...")
scanned_devices = self.perform_network_scan()
diff = self.update_cache_and_diff(scanned_devices)
if self._on_diff is not None:
try:
self._on_diff(diff)
except Exception: # noqa: BLE001
logger.exception("on_diff hook failed")
significant = bool(diff.get("new"))
raw_message = self.produce_notification(diff)
now = time.time()
def flush_minor():
nonlocal pending, pending_first_ts
if pending:
self.send_notification("\n\n".join(pending), notif_type="nmap", disable_notif=True)
pending = []
pending_first_ts = None
if raw_message:
if significant:
# Важное событие: сначала отправляем накопленные мелочи, затем само событие со звуком.
flush_minor()
self.send_notification(raw_message, notif_type="nmap", disable_notif=False)
else:
pending.append(raw_message)
if pending_first_ts is None:
pending_first_ts = now
combined_len = sum(len(m) for m in pending)
if (combined_len >= config.AGGREGATION_MAX_CHARS
or (now - pending_first_ts) >= config.AGGREGATION_INTERVAL):
flush_minor()
else:
if pending and (now - pending_first_ts) >= config.AGGREGATION_INTERVAL:
flush_minor()
else:
with self._queue_lock:
nmap_queue = [msg for msg in self.load_message_queue() if msg.get("type") == "nmap"]
if nmap_queue:
self.send_notification("", notif_type="nmap", disable_notif=True)
time.sleep(self.long_check_interval)
def start_scan(self):
"""Долгие проверки (сканирование сети). Мониторы доп. источников
планирует модуль network (register_monitor), а не сканер."""
Thread(target=self.start_long_checks).start()
store.py 6.1 КБ
"""network.store — единственное место IO кеша сканера и игнор-листа.
Владелец файлов networkscan_cache.json / networkscan_ignored.json — пакет
network: сканер пишет, ui_web и admin-движок (nmap-страницы TG) читают/правят
только через эти функции, а не прямым open().
"""
import json
import os
def _paths():
from . import config as pkg_config
return pkg_config.CACHE_FILE_NETWORK, pkg_config.SCAN_NETWORK_IGNORED_MACS
def _read_json(path):
try:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return None
def _write_json(path, data):
tmp = path + ".tmp"
os.makedirs(os.path.dirname(tmp) or ".", exist_ok=True)
with open(tmp, "w", encoding="utf-8") as f:
json.dump(data, f, indent=4, ensure_ascii=False)
os.replace(tmp, path)
def load_devices():
cache_file, _ = _paths()
data = _read_json(cache_file)
return data if isinstance(data, dict) else {}
def save_devices(cache) -> None:
cache_file, _ = _paths()
_write_json(cache_file, cache or {})
def set_device_name(mac, name):
"""Кастомное имя устройства. Вернуть (ok, err)."""
cache = load_devices()
if mac not in cache:
return False, "MAC не найден в кэше"
cache[mac]["name_custom"] = name
try:
save_devices(cache)
except Exception as e: # noqa: BLE001
return False, str(e)
return True, ""
def delete_device(mac):
"""Удалить запись устройства из кеша. Вернуть (ok, err).
Устройство вернётся при следующем скане, если оно реально в сети, —
это способ сбросить историю/стат, а не вечный бан (для бана —
игнор-лист)."""
cache = load_devices()
if mac not in cache:
return False, "MAC не найден в кэше"
try:
del cache[mac]
save_devices(cache)
except Exception as e: # noqa: BLE001
return False, str(e)
return True, ""
def load_ignored():
_, ignored_file = _paths()
data = _read_json(ignored_file)
return data if isinstance(data, list) else []
def devices_snapshot():
"""Обогащённый снимок устройств для панели мониторинга.
Порт агрегата из transport_web.monitoring (переехал к владельцу данных):
durations, сортировка онлайн-first. Best-effort — пустой кеш = пусто.
"""
import time as _time
def _fmt_ts(ts):
if not ts:
return None
return _time.strftime("%Y-%m-%d %H:%M:%S", _time.localtime(ts))
def _fmt_duration(sec):
if sec is None or sec < 0:
return "—"
sec = int(sec)
if sec < 60:
return f"{sec} сек"
m, s = divmod(sec, 60)
h, m = divmod(m, 60)
d, h = divmod(h, 24)
parts = []
if d:
parts.append(f"{d} д")
if h:
parts.append(f"{h} ч")
if m:
parts.append(f"{m} мин")
if not parts and s:
parts.append(f"{s} сек")
return " ".join(parts) or "0 мин"
scan = load_devices()
now_ts = _time.time()
devices = []
for mac, d in scan.items():
connected = bool(d.get("connected"))
online_start = d.get("online_start_ts")
last_appearance = d.get("last_appearance")
last_disconnection = d.get("last_disconnection")
first_appearance = d.get("first_appearance")
if connected:
last_change_ts = online_start if online_start else (last_appearance or first_appearance)
else:
last_change_ts = last_disconnection if last_disconnection else (last_appearance or first_appearance)
duration_s = int(now_ts - last_change_ts) if last_change_ts else 0
if duration_s < 0: # защита от будущих меток в кеше
duration_s = 0
devices.append({
"name": d.get("name_custom") or d.get("name") or mac,
"ip": d.get("ip"),
"mac": mac,
"online": connected,
"ports": d.get("ports") or [],
"last_seen": _fmt_ts(last_appearance),
"offline_since": _fmt_ts(last_disconnection),
"first_appearance": _fmt_ts(first_appearance),
"last_change": _fmt_ts(last_change_ts),
"last_change_ts": last_change_ts,
"duration_s": duration_s,
"duration_str": _fmt_duration(duration_s) if last_change_ts else "—",
"connections_total": d.get("connections_total"),
"connections_today": d.get("connections_today"),
"online_seconds_total": d.get("online_seconds_total"),
"online_seconds_today": d.get("online_seconds_today"),
})
devices.sort(key=lambda x: (not x["online"], (x["name"] or "").lower()))
return {"total": len(devices),
"online": sum(1 for d in devices if d["online"]),
"devices": devices}
def save_ignored(ignored) -> None:
_, ignored_file = _paths()
_write_json(ignored_file, ignored or [])
def set_ignored(action, mac):
"""add/remove MAC в игнор-листе. Вернуть (ignored|None, err)."""
ignored = load_ignored()
ignored_up = [x.upper() for x in ignored]
if action == "add":
if mac in ignored_up:
return None, "MAC уже в игнор-листе"
if mac not in load_devices():
return None, "MAC не найден в кэше сети"
ignored.append(mac)
else:
if mac not in ignored_up:
return None, "MAC нет в игнор-листе"
ignored = [x for x in ignored if x.upper() != mac]
try:
save_ignored(ignored)
except Exception as e: # noqa: BLE001
return None, str(e)
return ignored, ""
tg.py 16.3 КБ
"""network.tg — TG-меню сканирования (паритет с продом).
Статус/статистика с пагинацией, кастомные имена MAC, игнор-лист.
Данные — devices_snapshot + store пакета (имена/игнор).
Просмотр — view_net_status, изменения — nmap_manage.
Legacy-регистрация в ui_tg (группа -1; ввод текста — группа -6).
"""
import logging
import re
from telegram import InlineKeyboardButton, InlineKeyboardMarkup
from telegram.ext import DispatcherHandlerStop
from . import store as _store
logger = logging.getLogger(__name__)
_VIEW_RIGHT = "view_net_status"
_EDIT_RIGHT = "nmap_manage"
_PAGE = 8
_MAC_RE = r"^([0-9A-Fa-f]{2}[:-]){5}[0-9A-Fa-f]{2}$"
def _qid(uid):
s = str(uid)
return s if ":" in s else "telegram:%s" % s
def _display(d):
return (d.get("name_custom") or d.get("name") or d.get("mac")
or "Unknown")
class NetMenu:
def __init__(self, ctx):
self.ctx = ctx
def _can_view(self, user_id):
try:
return bool(self.ctx.rights.can(_qid(user_id), _VIEW_RIGHT))
except Exception: # noqa: BLE001
return False
def _can_edit(self, user_id):
try:
return bool(self.ctx.rights.can(_qid(user_id), _EDIT_RIGHT))
except Exception: # noqa: BLE001
return False
def _devices(self):
try:
snap = self.ctx.api.call("network.devices_snapshot") or {}
except Exception as e: # noqa: BLE001
logger.error("network TG snapshot: %s", e)
return []
devs = snap.get("devices") or []
return [d for d in devs if isinstance(d, dict)]
def _send(self, bot, chat_id, text, reply_markup=None):
try:
bot.send_message(chat_id=chat_id, text=text,
reply_markup=reply_markup)
except Exception as e: # noqa: BLE001
logger.error("network TG send: %s", e)
# -- меню ---------------------------------------------------------
def show_menu(self, update, context):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_view(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
kb = InlineKeyboardMarkup([
[InlineKeyboardButton("📊 Статус сети",
callback_data="net_status")],
[InlineKeyboardButton("📈 Статистика устройств",
callback_data="net_stats")],
[InlineKeyboardButton("📝 Изменить имя для MAC",
callback_data="net_addname")],
[InlineKeyboardButton("🔇 Фильтр MAC",
callback_data="net_ignore")],
[InlineKeyboardButton("🏠 Главное меню",
callback_data="main_menu")],
])
self._send(context.bot, chat_id, "Меню сканирования:", kb)
# -- статус --------------------------------------------------------
def show_status(self, update, context, page=1):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_view(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
devs = self._devices()
if not devs:
self._send(context.bot, chat_id, "Ошибка загрузки кэша сети.")
return
total_pages = max(1, (len(devs) + _PAGE - 1) // _PAGE)
page = max(1, min(page, total_pages))
chunk = devs[(page - 1) * _PAGE:page * _PAGE]
header = "Обзор сети"
if total_pages > 1:
header += " (стр. %s/%s)" % (page, total_pages)
lines = [header, ""]
for d in chunk:
st = "🔄" if d.get("online") else "🔌"
lines.append("%s %s (%s) - %s" % (
st, d.get("mac") or "?", _display(d),
d.get("ip") or "Unknown"))
lines += ["", "Выберете MAC для изменения имени."]
mac_buttons = []
for d in chunk:
mac = d.get("mac") or ""
if not mac:
continue
mac_buttons.append(InlineKeyboardButton(
_display(d), callback_data="net_setname:%s" % mac))
rows = [mac_buttons[i:i + 3] for i in range(0, len(mac_buttons), 3)]
nav = []
if page > 1:
nav.append(InlineKeyboardButton(
"⬅️", callback_data="net_status:%s" % (page - 1)))
nav.append(InlineKeyboardButton("%s/%s" % (page, total_pages),
callback_data="net_noop"))
if page < total_pages:
nav.append(InlineKeyboardButton(
"➡️", callback_data="net_status:%s" % (page + 1)))
rows.append(nav)
rows.append([InlineKeyboardButton("🔙 Назад",
callback_data="net_menu")])
self._send(context.bot, chat_id, "\n".join(lines),
InlineKeyboardMarkup(rows))
# -- статистика ------------------------------------------------------
def show_stats(self, update, context, page=1):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_view(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
devs = self._devices()
if not devs:
self._send(context.bot, chat_id, "Ошибка загрузки кэша сети.")
return
online = sum(1 for d in devs if d.get("online"))
text = ("Всего устройств: %s · Онлайн: %s · Офлайн: %s\n\n" % (
len(devs), online, len(devs) - online))
total_pages = max(1, (len(devs) + _PAGE - 1) // _PAGE)
page = max(1, min(page, total_pages))
chunk = devs[(page - 1) * _PAGE:page * _PAGE]
lines = [text]
for d in chunk:
st = "🔄" if d.get("online") else "🔌"
lines.append("%s %s (%s)" % (st, _display(d),
d.get("ip") or "Unknown"))
lines.append(" коннектов: %s/%s · онлайн: %s" % (
d.get("connections_today", 0),
d.get("connections_total", 0),
_fmt_online(d.get("online_seconds_total", 0))))
if total_pages > 1:
lines.append("(стр. %s/%s)" % (page, total_pages))
nav = []
if page > 1:
nav.append(InlineKeyboardButton(
"⬅️", callback_data="net_stats:%s" % (page - 1)))
if page < total_pages:
nav.append(InlineKeyboardButton(
"➡️", callback_data="net_stats:%s" % (page + 1)))
rows = [nav] if nav else []
rows.append([InlineKeyboardButton("🔙 Назад",
callback_data="net_menu")])
self._send(context.bot, chat_id, "\n".join(lines),
InlineKeyboardMarkup(rows))
# -- переименование ----------------------------------------------------
def ask_mac(self, update, context):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_edit(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
context.chat_data["awaiting"] = "net_addname"
self._send(context.bot, chat_id,
"Введите MAC адрес для которого хотите изменить имя:")
def ask_name(self, update, context, mac):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_edit(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
context.chat_data["awaiting"] = "net_setname"
context.chat_data["mac_for_setname"] = mac
self._send(context.bot, chat_id,
"Введите новое имя для MAC %s:" % mac)
# -- игнор -------------------------------------------------------------
def show_ignore(self, update, context):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_edit(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
try:
ignored = _store.load_ignored() or []
except Exception: # noqa: BLE001
ignored = []
text = ("🔇 Игнор-лист (%s):\n%s" % (
len(ignored),
"\n".join("• %s" % m for m in ignored) if ignored else "— пусто —"))
kb = InlineKeyboardMarkup([
[InlineKeyboardButton("➕ Добавить MAC",
callback_data="net_ignore_add")],
[InlineKeyboardButton("❌ Удалить MAC",
callback_data="net_ignore_del")],
[InlineKeyboardButton("🔙 Назад", callback_data="net_menu")],
])
self._send(context.bot, chat_id, text, kb)
def ask_ignore(self, update, context, action):
query = update.callback_query
try:
query.answer()
except Exception: # noqa: BLE001
pass
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_edit(user_id):
self._send(context.bot, chat_id, "🚫 Доступ запрещён!")
return
context.chat_data["awaiting"] = ("net_ignore_add"
if action == "add"
else "net_ignore_del")
self._send(context.bot, chat_id,
"Введите MAC адрес для %s в игнор-лист:" % (
"добавления" if action == "add" else "удаления"))
# -- ввод текстом --------------------------------------------------------
def on_text(self, update, context):
state = context.chat_data.get("awaiting")
if state not in ("net_addname", "net_setname", "net_ignore_add",
"net_ignore_del"):
return
context.chat_data.pop("awaiting", None)
try:
user_id = update.effective_user.id
chat_id = update.effective_chat.id
if not self._can_edit(user_id):
return
text = (update.message.text or "").strip()
if state == "net_addname":
if not re.match(_MAC_RE, text):
self._send(context.bot, chat_id,
"❌ Некорректный MAC. Формат: AA:BB:CC:DD:EE:FF.")
return
context.chat_data["awaiting"] = "net_setname"
context.chat_data["mac_for_setname"] = text.upper()
self._send(context.bot, chat_id,
"Введите новое имя для MAC %s:" % text.upper())
return
if state == "net_setname":
mac = context.chat_data.pop("mac_for_setname", "")
if not mac or not text:
self._send(context.bot, chat_id, "❌ Отменено.")
return
ok, err = _store.set_device_name(mac, text[:64])
self._send(context.bot, chat_id,
"✅ Имя сохранено: %s" % mac if ok
else "⚠️ %s" % (err or "ошибка"))
return
action = "add" if state == "net_ignore_add" else "remove"
if not re.match(_MAC_RE, text):
self._send(context.bot, chat_id,
"❌ Некорректный MAC. Формат: AA:BB:CC:DD:EE:FF.")
return
_, err = _store.set_ignored(action, text.upper())
self._send(context.bot, chat_id,
"✅ Игнор-лист обновлён." if not err
else "⚠️ %s" % err)
finally:
# Не пускаем ввод в общий текстовый хендлер admin.
raise DispatcherHandlerStop()
# -- вход ------------------------------------------------------------------
def handle(self, update, context):
data = (update.callback_query.data or "")
# Правило меню: предыдущее затирается, новое отправляется.
# main_menu — дальше admin покажет главное (без стопа).
try:
context.bot.delete_message(
update.callback_query.message.chat_id,
update.callback_query.message.message_id)
except Exception: # noqa: BLE001
pass
if data == "main_menu":
return
if data == "net_noop":
try:
update.callback_query.answer()
except Exception: # noqa: BLE001
pass
return
if data in ("net_menu", "net_status"):
if data == "net_menu":
self.show_menu(update, context)
else:
self.show_status(update, context, 1)
elif data.startswith("net_status:"):
try:
page = int(data.split(":", 1)[1])
except ValueError:
page = 1
self.show_status(update, context, page)
elif data in ("net_stats",):
self.show_stats(update, context, 1)
elif data.startswith("net_stats:"):
try:
page = int(data.split(":", 1)[1])
except ValueError:
page = 1
self.show_stats(update, context, page)
elif data == "net_addname":
self.ask_mac(update, context)
elif data.startswith("net_setname:"):
self.ask_name(update, context, data.split(":", 1)[1])
elif data.startswith("nmap_setname:"):
# Легаси-кнопки из уведомлений сканера (scan.py): тот же
# диалог, что net_setname. Без этой ветки — «Неизвестная
# команда» из catch-all admin.
self.ask_name(update, context, data.split(":", 1)[1])
elif data == "net_ignore":
self.show_ignore(update, context)
elif data == "net_ignore_add":
self.ask_ignore(update, context, "add")
elif data == "net_ignore_del":
self.ask_ignore(update, context, "del")
# Не пускаем свои колбэки в catch-all admin.button (иначе
# «Неизвестная команда» поверх). main_menu — выше без стопа.
raise DispatcherHandlerStop()
def _fmt_online(sec):
try:
sec = int(sec or 0)
except (TypeError, ValueError):
return "—"
if sec < 60:
return "%s сек" % sec
m, s = divmod(sec, 60)
h, m = divmod(m, 60)
d, h = divmod(h, 24)
parts = []
if d:
parts.append("%s д" % d)
if h:
parts.append("%s ч" % h)
if m:
parts.append("%s мин" % m)
return " ".join(parts) or "0 мин"
ui_tg.py 1.2 КБ
"""network.ui_tg — кнопка /start + legacy-хендлер меню сканирования."""
TG_COMMANDS = {}
TG_MENU = [
{"id": "netscan", "title": "📡 Сканирование", "callback": "net_menu",
"right": "view_net_status"},
]
def register_ptb(facade, mod):
"""Legacy-хук: меню сети на фасаде напрямую (колбэки net_*).
Группа -1: раньше catch-all admin.button (группа 0).
Ввод текста — группа -5 (НЕ -6: в PTB внутри группы срабатывает
только первый совпавший хендлер; в -6 уже сидит fail2ban.on_text,
и network.on_text за ним не вызывался бы никогда).
"""
from telegram.ext import CallbackQueryHandler, MessageHandler, Filters
from .tg import NetMenu
menu = NetMenu(mod.ctx)
mod.net_menu = menu
facade.add_handler(CallbackQueryHandler(
menu.handle,
pattern=(r"^(net_menu|net_status|net_stats|net_addname|"
r"net_setname:|net_ignore|net_noop|"
r"main_menu|nmap_setname:)")), group=-1)
facade.add_handler(MessageHandler(Filters.text, menu.on_text),
group=-5)
ui_web.py 4.5 КБ
"""network.ui_web — /api/network/devices + /api/admin/nmap*.
Порт AdminNmap/Name/Ignore из старого messengers/web.py. ACL и traffic
переехали в модуль network_router. Пути кешей — через ctx.api network.
"""
ROUTES = [
("GET", "/api/network/devices", "net_devices", {"right": None,
"desc": "Устройства сети для панели мониторинга (свои кеши)"}),
("GET", "/api/admin/nmap", "nmap", {"right": "nmap_manage",
"desc": "Устройства сканера + игнор-лист"}),
("POST", "/api/admin/nmap/name", "nmap_name", {"right": "nmap_manage",
"desc": "Кастомное имя {mac, name}"}),
("POST", "/api/admin/nmap/ignore", "nmap_ignore", {"right": "nmap_manage",
"desc": "Игнор-лист {action: add|remove, mac}"}),
("POST", "/api/network/nmap/delete", "nmap_delete", {"right": "nmap_manage",
"desc": "Удалить устройство из кеша {mac}"}),
]
SERVICE = [
{"key": "nmap_manage", "icon": "🔍", "title": "Сканер сети", "right": "nmap_manage",
"phase": 3, "desc": "Устройства, кастомные имена, игнор-лист"},
]
PANELS = []
MONITOR_PANELS = [
{"id": "net.devices", "title": "Сеть: устройства", "api": "/api/network/devices",
"visibility": "public", "refresh_s": 60},
]
_MAC_RE = r"^([0-9A-Fa-f]{2}[:-]){5}[0-9A-Fa-f]{2}$"
from . import store as _store
def handle_api(ctx, config, method, req):
import re as _re
from core.errors import UserError
if method == "net_devices":
try:
data = ctx.api.call("network.devices_snapshot")
except Exception as e: # noqa: BLE001
return 500, {"ok": False, "error": str(e)[:200]}
if not isinstance(data, dict):
data = {"total": 0, "online": 0, "devices": []}
data["ok"] = True
return data
from core.errors import UserError
if method == "nmap":
cache = _store.load_devices()
ignored = _store.load_ignored()
devices = []
for mac, d in cache.items():
devices.append({
"mac": mac, "ip": d.get("ip"), "name": d.get("name") or "",
"name_custom": d.get("name_custom") or "",
"type": d.get("type") or "", "online": bool(d.get("connected")),
"ports": d.get("ports") or [],
"ignored": mac.upper() in [x.upper() for x in ignored],
})
devices.sort(key=lambda x: (not x["online"],
(x["name_custom"] or x["name"] or "").lower()))
return {"devices": devices, "ignored": ignored}
if method == "nmap_name":
body = req.get("body") or {}
mac = (body.get("mac") or "").strip().upper()
name = (body.get("name") or "").strip()
if not _re.match(_MAC_RE, mac):
return 400, {"ok": False, "error": "Некорректный MAC"}
if len(name) > 64:
return 400, {"ok": False, "error": "Имя слишком длинное (64)"}
ok, err = _store.set_device_name(mac, name)
if not ok:
code = 400 if "не найден" in err else 500
return code, {"ok": False, "error": err[:400]}
return {"ok": True, "mac": mac, "name": name}
if method == "nmap_ignore":
body = req.get("body") or {}
action = (body.get("action") or "").strip()
mac = (body.get("mac") or "").strip().upper()
if action not in ("add", "remove"):
return 400, {"ok": False, "error": "action add/remove"}
if not _re.match(_MAC_RE, mac):
return 400, {"ok": False, "error": "Некорректный MAC"}
ignored, err = _store.set_ignored(action, mac)
if err:
code = 400 if "MAC " in err or "уже" in err or "нет" in err else 500
return code, {"ok": False, "error": err[:400]}
return {"ok": True, "ignored": ignored}
if method == "nmap_delete":
body = req.get("body") or {}
mac = (body.get("mac") or "").strip().upper()
if not _re.match(_MAC_RE, mac):
return 400, {"ok": False, "error": "Некорректный MAC"}
ok, err = _store.delete_device(mac)
if not ok:
code = 400 if "не найден" in err else 500
return code, {"ok": False, "error": err[:400]}
return {"ok": True, "mac": mac}
raise UserError("неизвестный метод: %s" % method)