From dc161e151d20c53f0d897a4481bd69d2d0960187 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=90=D0=BD=D0=B4=D1=80=D0=B5=D0=B9=20=D0=91=D0=BE=D0=B1?= =?UTF-8?q?=D1=8B=D1=80=D0=B5=D0=B2?= Date: Fri, 22 May 2026 03:05:34 +0300 Subject: [PATCH] fix(keenetic): try all DNS A records on auth failure KeenDNS returns multiple A records including unreachable IPs; aiohttp was timing out on the first bad address. Iterate all resolved IPs with short connect timeouts, serialize polls with a lock, and fail faster. Co-authored-by: Cursor --- server/api/keenetic.py | 54 ++++---- server/services/keenetic_client.py | 211 +++++++++++++++++++---------- 2 files changed, 170 insertions(+), 95 deletions(-) diff --git a/server/api/keenetic.py b/server/api/keenetic.py index 840b926..6ecfa82 100644 --- a/server/api/keenetic.py +++ b/server/api/keenetic.py @@ -23,9 +23,10 @@ logger = logging.getLogger(__name__) keenetic_metrics: Dict[str, dict] = {} _refresh_all_running = False +_poll_lock = asyncio.Lock() KEENETIC_FILE = DATA_DIR / "keenetic.json" -DEVICE_REFRESH_TIMEOUT = 75 +DEVICE_REFRESH_TIMEOUT = 45 REFRESH_ALL_GAP_SEC = 2 @@ -95,6 +96,11 @@ async def _refresh_device(dev: dict) -> dict: await client.close() +async def _refresh_device_locked(dev: dict) -> dict: + async with _poll_lock: + return await _refresh_device(dev) + + @router.get("/list") async def keenetic_list(request: Request, user: str = Depends(require_auth)): devices = _load_keenetic() @@ -264,7 +270,7 @@ async def keenetic_update(name: str, request: Request, user: str = Depends(requi result = {"status": "ok", "web_url": web_url, "host": host} if body.get("refresh", True): try: - metrics = await _refresh_device(devices[idx]) + metrics = await _refresh_device_locked(devices[idx]) result["metrics"] = metrics except Exception as e: result["refresh_error"] = str(e) @@ -289,7 +295,7 @@ async def keenetic_refresh(name: str, request: Request, user: str = Depends(requ return {"status": "error", "detail": "router not found"} try: - metrics = await _refresh_device(dev) + metrics = await _refresh_device_locked(dev) if not metrics["online"] and metrics.get("error"): return {"status": "error", "detail": metrics["error"], "metrics": metrics} return {"status": "ok", "metrics": metrics} @@ -302,18 +308,19 @@ async def _refresh_all_devices() -> list: devices = _load_keenetic() results = [] try: - for i, dev in enumerate(devices): - try: - metrics = await _refresh_device(dev) - results.append({ - "name": dev["name"], - "online": metrics["online"], - "error": metrics.get("error", ""), - }) - except Exception as e: - results.append({"name": dev["name"], "online": False, "error": str(e)}) - if i + 1 < len(devices): - await asyncio.sleep(REFRESH_ALL_GAP_SEC) + async with _poll_lock: + for i, dev in enumerate(devices): + try: + metrics = await _refresh_device(dev) + results.append({ + "name": dev["name"], + "online": metrics["online"], + "error": metrics.get("error", ""), + }) + except Exception as e: + results.append({"name": dev["name"], "online": False, "error": str(e)}) + if i + 1 < len(devices): + await asyncio.sleep(REFRESH_ALL_GAP_SEC) finally: _refresh_all_running = False return results @@ -374,18 +381,19 @@ async def keenetic_reboot(name: str, request: Request, user: str = Depends(requi async def keenetic_monitor_loop(): """Background polling for all Keenetic routers.""" - await asyncio.sleep(15) + await asyncio.sleep(5) while True: try: devices = _load_keenetic() - if devices: + if devices and not _refresh_all_running: logger.info(f"Keenetic monitor: refreshing {len(devices)} routers") - for dev in devices: - try: - await _refresh_device(dev) - except Exception as e: - logger.error(f"Keenetic refresh {dev['name']}: {e}") - await asyncio.sleep(REFRESH_ALL_GAP_SEC) + async with _poll_lock: + for dev in devices: + try: + await _refresh_device(dev) + except Exception as e: + logger.error(f"Keenetic refresh {dev['name']}: {e}") + await asyncio.sleep(REFRESH_ALL_GAP_SEC) except Exception as e: logger.error(f"Keenetic monitor loop error: {e}") diff --git a/server/services/keenetic_client.py b/server/services/keenetic_client.py index 5c494aa..f0fcccb 100644 --- a/server/services/keenetic_client.py +++ b/server/services/keenetic_client.py @@ -13,7 +13,8 @@ import aiohttp logger = logging.getLogger(__name__) KEENDNS_MARKERS = (".pro", ".club", ".link", "netcraze", "keenetic") -AUTH_RETRIES = 3 +AUTH_RETRIES = 2 +BLOCKED_IPS = frozenset({"0.0.0.0", "127.0.0.1"}) _IP_HOST_RE = re.compile(r"^\d{1,3}(?:\.\d{1,3}){3}$") @@ -56,11 +57,34 @@ def is_keendns_host(host: str) -> bool: return any(m in domain for m in KEENDNS_MARKERS) -def client_timeout_for(host: str) -> aiohttp.ClientTimeout: +def client_timeout_for(host: str, *, probing: bool = False) -> aiohttp.ClientTimeout: """Timeouts tuned for KeenDNS vs direct IP.""" + if probing: + return aiohttp.ClientTimeout(total=15, connect=4, sock_read=10) if is_public_ip_host(host): return aiohttp.ClientTimeout(total=20, connect=8, sock_read=12) - return aiohttp.ClientTimeout(total=60, connect=20, sock_read=40) + return aiohttp.ClientTimeout(total=45, connect=12, sock_read=30) + + +async def resolve_ipv4_addresses(hostname: str) -> list[str]: + """Resolve KeenDNS A records, skipping placeholder/blocked addresses.""" + if not hostname or is_public_ip_host(hostname.split(":")[0]): + return [] + loop = asyncio.get_running_loop() + try: + infos = await loop.getaddrinfo( + hostname, None, family=socket.AF_INET, type=socket.SOCK_STREAM, + ) + except socket.gaierror: + return [] + result, seen = [], set() + for info in infos: + ip = info[4][0] + if ip in BLOCKED_IPS or ip in seen: + continue + seen.add(ip) + result.append(ip) + return result def build_api_base_url(host: str, web_url: str = "") -> str: @@ -87,6 +111,7 @@ def _make_connector() -> aiohttp.TCPConnector: family=socket.AF_INET, force_close=True, enable_cleanup_closed=True, + ttl_dns_cache=30, ) @@ -108,22 +133,47 @@ class KeeneticClient: self.base_url = base parsed = urlparse(base) self._host_key = parsed.netloc or host + self._hostname = parsed.hostname or "" + self._port = parsed.port or (443 if parsed.scheme == "https" else 80) + self._scheme = parsed.scheme or "https" + self._netloc = parsed.netloc or host + self._active_target: Optional[str] = None + self._resolved_ips: list[str] = [] 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: + def _effective_base(self) -> str: + if self._active_target: + port = f":{self._port}" if self._port not in (80, 443) else "" + return f"{self._scheme}://{self._active_target}{port}" + return self.base_url.rstrip("/") + + def _request_headers(self) -> dict: + headers = {"User-Agent": "VPS-Monitoring/1.0"} + if self._active_target: + headers["Host"] = self._netloc + return headers + + async def _connection_targets(self) -> list[str]: + if is_public_ip_host(self._hostname): + return [self._hostname] + if not self._resolved_ips: + self._resolved_ips = await resolve_ipv4_addresses(self._hostname) + return self._resolved_ips + [self._hostname] + + async def _get_session(self, *, probing: bool = False) -> aiohttp.ClientSession: if self._session is None or self._session.closed: - timeout = client_timeout_for(self._host_key) + timeout = client_timeout_for(self._host_key, probing=probing) 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"}, + headers=self._request_headers(), ) return self._session @@ -142,86 +192,103 @@ class KeeneticClient: return "Connection timeout" async def _authenticate_once(self) -> bool: - """Single auth attempt.""" + """Single auth attempt; tries each KeenDNS A record before giving up.""" 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): + targets = await self._connection_targets() + last_err: Optional[Exception] = None + + for target in targets: + await self._reset_session() + self._active_target = target if is_public_ip_host(target) else None + auth_url = f"{self._effective_base()}/auth" + try: - return await self._authenticate_once() + session = await self._get_session(probing=True) + 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} @ {auth_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() + + session = await self._get_session() + 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} @ {auth_url}") + return False except aiohttp.ClientConnectorError as e: + last_err = 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 + return False + continue 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__}" + last_err = e + logger.debug( + f"Keenetic auth try failed @ {auth_url}: {type(e).__name__}" ) + continue + + if last_err: + if isinstance(last_err, aiohttp.ClientConnectorError): + self.last_error = "Cannot connect to router" + else: + self.last_error = self._timeout_error_message() + else: + self.last_error = self._timeout_error_message() + return False + + async def authenticate(self) -> bool: + """Perform challenge-response authentication with retries.""" + for attempt in range(AUTH_RETRIES): + if attempt: + self._resolved_ips = await resolve_ipv4_addresses(self._hostname) await self._reset_session() - if attempt + 1 < AUTH_RETRIES: - await asyncio.sleep(2.0 * (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}") + if await self._authenticate_once(): + return True + if self.last_error == "DNS не резолвится с VPS": return False + if attempt + 1 < AUTH_RETRIES: + logger.warning( + f"Keenetic auth retry ({attempt + 1}/{AUTH_RETRIES}) " + f"@ {self.base_url}: {self.last_error}" + ) + await asyncio.sleep(1.5 * (attempt + 1)) return False async def rci_show(self, command: str, params: Optional[dict] = None) -> Optional[dict]: @@ -232,7 +299,7 @@ class KeeneticClient: session = await self._get_session() path = command.replace(" ", "/") - url = f"{self.base_url}/rci/show/{path}" + url = f"{self._effective_base()}/rci/show/{path}" try: async with session.get(url, params=params, ssl=False) as resp: @@ -262,7 +329,7 @@ class KeeneticClient: asyncio.TimeoutError, TimeoutError, ) - url = f"{self.base_url}/rci/" + url = f"{self._effective_base()}/rci/" for attempt in range(AUTH_RETRIES): try: session = await self._get_session()