mirror of
https://github.com/andrey271192/vps_monitoring.git
synced 2026-09-20 11:55:34 +00:00
Use web_url scheme and host:port for API base URL instead of forcing https. Force IPv4 in aiohttp to avoid IPv6 hangs on KeenDNS. Sync canonical 24-router list with Russian names. Co-authored-by: Cursor <cursoragent@cursor.com>
451 lines
17 KiB
Python
451 lines
17 KiB
Python
"""Keenetic router RCI API client for monitoring via KeenDNS."""
|
||
|
||
import asyncio
|
||
import hashlib
|
||
import logging
|
||
import re
|
||
import socket
|
||
from typing import Optional
|
||
from urllib.parse import urlparse
|
||
|
||
import aiohttp
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
KEENDNS_MARKERS = (".pro", ".club", ".link", "netcraze", "keenetic")
|
||
AUTH_RETRIES = 3
|
||
_IP_HOST_RE = re.compile(r"^\d{1,3}(?:\.\d{1,3}){3}$")
|
||
|
||
|
||
def normalize_web_url(url: str) -> str:
|
||
"""Ensure web UI URL has a scheme (https by default)."""
|
||
url = (url or "").strip()
|
||
if not url:
|
||
return ""
|
||
if not url.startswith("http://") and not url.startswith("https://"):
|
||
url = "https://" + url
|
||
return url.rstrip("/")
|
||
|
||
|
||
def is_public_ip_host(host: str) -> bool:
|
||
"""True when host is a bare IPv4 (optional :port), not KeenDNS."""
|
||
if not host:
|
||
return False
|
||
return bool(_IP_HOST_RE.match(host.split(":")[0]))
|
||
|
||
|
||
def is_keendns_host(host: str) -> bool:
|
||
domain = (host or "").split(":")[0].lower()
|
||
return any(m in domain for m in KEENDNS_MARKERS)
|
||
|
||
|
||
def client_timeout_for(host: str) -> aiohttp.ClientTimeout:
|
||
"""Timeouts tuned for KeenDNS vs direct IP."""
|
||
if is_public_ip_host(host):
|
||
return aiohttp.ClientTimeout(total=20, connect=8, sock_read=12)
|
||
return aiohttp.ClientTimeout(total=45, connect=15, sock_read=30)
|
||
|
||
|
||
def build_api_base_url(host: str, web_url: str = "") -> str:
|
||
"""Build RCI API base URL — prefer scheme/host/port from web_url."""
|
||
if web_url:
|
||
parsed = urlparse(normalize_web_url(web_url))
|
||
if parsed.scheme and parsed.netloc:
|
||
return f"{parsed.scheme}://{parsed.netloc}".rstrip("/")
|
||
|
||
raw = (host or "").strip()
|
||
if not raw:
|
||
return ""
|
||
|
||
if raw.startswith("http://") or raw.startswith("https://"):
|
||
return raw.rstrip("/")
|
||
|
||
scheme = "https" if is_keendns_host(raw) else "http"
|
||
return f"{scheme}://{raw}".rstrip("/")
|
||
|
||
|
||
def _make_connector() -> aiohttp.TCPConnector:
|
||
"""IPv4-only: avoids aiohttp hanging on broken IPv6 for KeenDNS multi-A records."""
|
||
return aiohttp.TCPConnector(
|
||
family=socket.AF_INET,
|
||
force_close=True,
|
||
enable_cleanup_closed=True,
|
||
)
|
||
|
||
|
||
class KeeneticClient:
|
||
"""Async client for Keenetic router RCI API.
|
||
|
||
Auth flow:
|
||
1. GET /auth -> 401 with X-NDM-Challenge + X-NDM-Realm headers
|
||
2. Compute: md5(login:realm:password) -> sha256(challenge + md5_hex)
|
||
3. POST /auth {"login": ..., "password": sha256_hex}
|
||
4. Session cookie persists for subsequent requests
|
||
"""
|
||
|
||
def __init__(self, host: str, login: str = "admin", password: str = "",
|
||
web_url: str = ""):
|
||
base = build_api_base_url(host, web_url)
|
||
if not base:
|
||
raise ValueError("host or web_url required")
|
||
self.base_url = base
|
||
parsed = urlparse(base)
|
||
self._host_key = parsed.netloc or host
|
||
self.login = login
|
||
self.password = password
|
||
self._session: Optional[aiohttp.ClientSession] = None
|
||
self._authenticated = False
|
||
self.last_error = ""
|
||
|
||
async def _get_session(self) -> aiohttp.ClientSession:
|
||
if self._session is None or self._session.closed:
|
||
timeout = client_timeout_for(self._host_key)
|
||
jar = aiohttp.CookieJar(unsafe=True)
|
||
self._session = aiohttp.ClientSession(
|
||
timeout=timeout,
|
||
cookie_jar=jar,
|
||
connector=_make_connector(),
|
||
version=aiohttp.HttpVersion11,
|
||
headers={"User-Agent": "VPS-Monitoring/1.0"},
|
||
)
|
||
return self._session
|
||
|
||
async def _reset_session(self):
|
||
if self._session and not self._session.closed:
|
||
await self._session.close()
|
||
self._session = None
|
||
self._authenticated = False
|
||
|
||
def _timeout_error_message(self) -> str:
|
||
if is_public_ip_host(self._host_key):
|
||
return (
|
||
"HTTP API не отвечает с VPS (прямой IP). "
|
||
"Нужен KeenDNS или удалённый доступ Keenetic."
|
||
)
|
||
return "Connection timeout"
|
||
|
||
async def _authenticate_once(self) -> bool:
|
||
"""Single auth attempt."""
|
||
self.last_error = ""
|
||
session = await self._get_session()
|
||
auth_url = f"{self.base_url}/auth"
|
||
|
||
async with session.get(auth_url, ssl=False) as resp:
|
||
if resp.status == 200:
|
||
self._authenticated = True
|
||
return True
|
||
|
||
if resp.status != 401:
|
||
if resp.status in (400, 403):
|
||
self.last_error = "Wrong protocol or port (try http/https)"
|
||
else:
|
||
self.last_error = f"Auth HTTP {resp.status}"
|
||
logger.error(
|
||
f"Keenetic auth unexpected status: {resp.status} @ {self.base_url}"
|
||
)
|
||
return False
|
||
|
||
challenge = resp.headers.get("X-NDM-Challenge", "")
|
||
realm = resp.headers.get("X-NDM-Realm", "")
|
||
|
||
if not challenge or not realm:
|
||
self.last_error = "No auth challenge from router"
|
||
logger.error("Keenetic auth: missing challenge/realm headers")
|
||
return False
|
||
|
||
md5_input = f"{self.login}:{realm}:{self.password}"
|
||
md5_hex = hashlib.md5(md5_input.encode("utf-8")).hexdigest()
|
||
sha_input = f"{challenge}{md5_hex}"
|
||
sha_hex = hashlib.sha256(sha_input.encode("utf-8")).hexdigest()
|
||
|
||
async with session.post(
|
||
auth_url,
|
||
json={"login": self.login, "password": sha_hex},
|
||
ssl=False,
|
||
) as resp:
|
||
if resp.status == 200:
|
||
self._authenticated = True
|
||
return True
|
||
self.last_error = "Wrong login or password"
|
||
logger.error(f"Keenetic auth failed: {resp.status} @ {self.base_url}")
|
||
return False
|
||
|
||
async def authenticate(self) -> bool:
|
||
"""Perform challenge-response authentication with retries."""
|
||
retryable = (
|
||
aiohttp.ServerTimeoutError,
|
||
aiohttp.ClientOSError,
|
||
asyncio.TimeoutError,
|
||
TimeoutError,
|
||
)
|
||
for attempt in range(AUTH_RETRIES):
|
||
try:
|
||
return await self._authenticate_once()
|
||
except aiohttp.ClientConnectorError as e:
|
||
err = str(e).lower()
|
||
if "name or service not known" in err or "nodename nor servname" in err:
|
||
self.last_error = "DNS не резолвится с VPS"
|
||
else:
|
||
self.last_error = "Cannot connect to router"
|
||
logger.error(f"Keenetic auth error: {type(e).__name__}: {e}")
|
||
return False
|
||
except retryable as e:
|
||
self.last_error = self._timeout_error_message()
|
||
logger.warning(
|
||
f"Keenetic auth timeout ({attempt + 1}/{AUTH_RETRIES}) "
|
||
f"@ {self.base_url}: {type(e).__name__}"
|
||
)
|
||
await self._reset_session()
|
||
if attempt + 1 < AUTH_RETRIES:
|
||
await asyncio.sleep(1.5 * (attempt + 1))
|
||
continue
|
||
logger.error(f"Keenetic auth timeout @ {self.base_url}")
|
||
return False
|
||
except Exception as e:
|
||
self.last_error = type(e).__name__
|
||
logger.error(f"Keenetic auth error: {type(e).__name__}: {e}")
|
||
return False
|
||
return False
|
||
|
||
async def rci_show(self, command: str, params: Optional[dict] = None) -> Optional[dict]:
|
||
"""GET /rci/show/<command> with optional query params."""
|
||
if not self._authenticated:
|
||
if not await self.authenticate():
|
||
return None
|
||
|
||
session = await self._get_session()
|
||
path = command.replace(" ", "/")
|
||
url = f"{self.base_url}/rci/show/{path}"
|
||
|
||
try:
|
||
async with session.get(url, params=params, ssl=False) as resp:
|
||
if resp.status == 200:
|
||
return await resp.json(content_type=None)
|
||
elif resp.status == 401:
|
||
self._authenticated = False
|
||
if await self.authenticate():
|
||
async with session.get(url, params=params, ssl=False) as resp2:
|
||
if resp2.status == 200:
|
||
return await resp2.json(content_type=None)
|
||
logger.error(f"Keenetic RCI {command}: HTTP {resp.status}")
|
||
return None
|
||
except Exception as e:
|
||
logger.error(f"Keenetic RCI error ({command}): {type(e).__name__}: {e}")
|
||
return None
|
||
|
||
async def rci_post(self, body: dict) -> Optional[dict]:
|
||
"""POST /rci/ with JSON body for batch commands."""
|
||
if not self._authenticated:
|
||
if not await self.authenticate():
|
||
return None
|
||
|
||
retryable = (
|
||
aiohttp.ServerTimeoutError,
|
||
aiohttp.ClientOSError,
|
||
asyncio.TimeoutError,
|
||
TimeoutError,
|
||
)
|
||
url = f"{self.base_url}/rci/"
|
||
for attempt in range(AUTH_RETRIES):
|
||
try:
|
||
session = await self._get_session()
|
||
async with session.post(url, json=body, ssl=False) as resp:
|
||
if resp.status == 200:
|
||
return await resp.json(content_type=None)
|
||
if resp.status == 401:
|
||
self._authenticated = False
|
||
if await self.authenticate():
|
||
continue
|
||
logger.error(f"Keenetic RCI POST: HTTP {resp.status}")
|
||
return None
|
||
except retryable as e:
|
||
logger.warning(
|
||
f"Keenetic RCI POST timeout ({attempt + 1}/{AUTH_RETRIES}) "
|
||
f"@ {self.base_url}: {type(e).__name__}"
|
||
)
|
||
await self._reset_session()
|
||
if not await self.authenticate():
|
||
return None
|
||
if attempt + 1 < AUTH_RETRIES:
|
||
await asyncio.sleep(1.5 * (attempt + 1))
|
||
continue
|
||
logger.error(f"Keenetic RCI POST error: {type(e).__name__}: {e}")
|
||
return None
|
||
except Exception as e:
|
||
logger.error(f"Keenetic RCI POST error: {type(e).__name__}: {e}")
|
||
return None
|
||
return None
|
||
|
||
async def collect_metrics(self, cached_info: Optional[dict] = None) -> dict:
|
||
"""Lightweight refresh: system + internet + VPN only."""
|
||
result = {
|
||
"online": False,
|
||
"hostname": "",
|
||
"model": "",
|
||
"firmware": "",
|
||
"cpuload": 0,
|
||
"memtotal": 0,
|
||
"memfree": 0,
|
||
"mem_percent": 0,
|
||
"uptime": 0,
|
||
"uptime_str": "",
|
||
"internet": False,
|
||
"gateway_accessible": False,
|
||
"dns_accessible": False,
|
||
"vpn": [],
|
||
"clients_count": 0,
|
||
"wifi_clients": 0,
|
||
"wired_clients": 0,
|
||
"error": "",
|
||
}
|
||
|
||
if cached_info:
|
||
result["model"] = cached_info.get("model", "")
|
||
result["firmware"] = cached_info.get("firmware", "")
|
||
|
||
if not await self.authenticate():
|
||
result["error"] = self.last_error or "Authentication failed"
|
||
return result
|
||
|
||
result["online"] = True
|
||
|
||
batch_cmd = {
|
||
"show": {
|
||
"system": {},
|
||
"internet": {"status": {}},
|
||
"interface": {},
|
||
}
|
||
}
|
||
|
||
need_info = not result.get("model")
|
||
if need_info:
|
||
batch_cmd["show"]["version"] = {}
|
||
batch_cmd["show"]["defaults"] = {}
|
||
|
||
batch = await self.rci_post(batch_cmd)
|
||
|
||
if not batch or not isinstance(batch, dict):
|
||
result["error"] = "No data from router"
|
||
result["online"] = False
|
||
return result
|
||
|
||
show = batch.get("show", batch)
|
||
|
||
sys_data = show.get("system", {})
|
||
if sys_data:
|
||
result["hostname"] = sys_data.get("hostname", "")
|
||
result["cpuload"] = int(sys_data.get("cpuload", 0) or 0)
|
||
result["memtotal"] = int(sys_data.get("memtotal", 0) or 0)
|
||
result["memfree"] = int(sys_data.get("memfree", 0) or 0)
|
||
|
||
memtotal = result["memtotal"]
|
||
memfree = result["memfree"]
|
||
if memtotal > 0:
|
||
result["mem_percent"] = round((1 - memfree / memtotal) * 100, 1)
|
||
|
||
uptime_sec = int(sys_data.get("uptime", 0) or 0)
|
||
result["uptime"] = uptime_sec
|
||
if uptime_sec:
|
||
days = uptime_sec // 86400
|
||
hours = (uptime_sec % 86400) // 3600
|
||
mins = (uptime_sec % 3600) // 60
|
||
result["uptime_str"] = f"{days}d {hours}h {mins}m"
|
||
|
||
if need_info:
|
||
ver_data = show.get("version", {})
|
||
if ver_data:
|
||
result["firmware"] = ver_data.get("title", ver_data.get("release", ""))
|
||
|
||
defaults = show.get("defaults", {})
|
||
if defaults:
|
||
product = defaults.get("product", "")
|
||
hw_id = defaults.get("ndmhwid", "")
|
||
result["model"] = f"{product} ({hw_id})" if product and hw_id else product or hw_id
|
||
|
||
inet_block = show.get("internet", {})
|
||
inet = inet_block.get("status", inet_block) if isinstance(inet_block, dict) else {}
|
||
if inet:
|
||
result["internet"] = bool(inet.get("internet"))
|
||
result["gateway_accessible"] = bool(inet.get("gateway-accessible"))
|
||
result["dns_accessible"] = bool(inet.get("dns-accessible"))
|
||
|
||
VPN_TYPES = {"Wireguard", "WireGuard", "OpenVPN", "PPTP", "L2TP", "SSTP", "EoIP", "IPsec"}
|
||
|
||
ifaces = show.get("interface", {})
|
||
if ifaces and isinstance(ifaces, dict):
|
||
for iface_name, iface_data in ifaces.items():
|
||
if not isinstance(iface_data, dict):
|
||
continue
|
||
itype = iface_data.get("type", "")
|
||
if itype in VPN_TYPES:
|
||
result["vpn"].append({
|
||
"name": iface_name,
|
||
"type": itype,
|
||
"state": iface_data.get("state", ""),
|
||
"description": iface_data.get("description", ""),
|
||
"address": iface_data.get("address", ""),
|
||
})
|
||
|
||
return result
|
||
|
||
async def collect_detail(self) -> Optional[dict]:
|
||
"""Heavy detail fetch: interfaces + connected clients."""
|
||
if not self._authenticated:
|
||
if not await self.authenticate():
|
||
return None
|
||
|
||
batch = await self.rci_post({
|
||
"show": {
|
||
"interface": {},
|
||
"ip": {"hotspot": {}},
|
||
}
|
||
})
|
||
|
||
if not batch or not isinstance(batch, dict):
|
||
return None
|
||
|
||
show = batch.get("show", batch)
|
||
detail = {"interfaces": [], "clients": []}
|
||
|
||
VPN_TYPES = {"Wireguard", "WireGuard", "OpenVPN", "PPTP", "L2TP", "SSTP", "EoIP", "IPsec"}
|
||
SHOW_TYPES = {"GigabitEthernet", "XGigabitEthernet", "WifiMaster", "AccessPoint",
|
||
"Bridge", "PPPoE"} | VPN_TYPES
|
||
|
||
ifaces = show.get("interface", {})
|
||
if ifaces and isinstance(ifaces, dict):
|
||
for iface_name, iface_data in ifaces.items():
|
||
if not isinstance(iface_data, dict):
|
||
continue
|
||
itype = iface_data.get("type", "")
|
||
if itype in SHOW_TYPES:
|
||
detail["interfaces"].append({
|
||
"id": iface_data.get("id", iface_name),
|
||
"type": itype,
|
||
"description": iface_data.get("description", ""),
|
||
"state": iface_data.get("state", ""),
|
||
"address": iface_data.get("address", ""),
|
||
"uptime": iface_data.get("uptime", 0),
|
||
})
|
||
|
||
ip_block = show.get("ip", {})
|
||
hotspot = ip_block.get("hotspot", ip_block) if isinstance(ip_block, dict) else {}
|
||
if hotspot:
|
||
hosts = hotspot.get("host", [])
|
||
if isinstance(hosts, list):
|
||
for h in hosts:
|
||
detail["clients"].append({
|
||
"name": h.get("name", h.get("hostname", "")),
|
||
"hostname": h.get("hostname", ""),
|
||
"ip": h.get("ip", ""),
|
||
"mac": h.get("mac", ""),
|
||
"active": h.get("active", False),
|
||
"speed": h.get("speed", 0),
|
||
"ssid": h.get("ssid", ""),
|
||
})
|
||
|
||
return detail
|
||
|
||
async def close(self):
|
||
if self._session and not self._session.closed:
|
||
await self._session.close()
|