mirror of
https://github.com/andrey271192/vps_monitoring.git
synced 2026-09-20 11:55:34 +00:00
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 <cursoragent@cursor.com>
This commit is contained in:
@@ -23,9 +23,10 @@ logger = logging.getLogger(__name__)
|
|||||||
|
|
||||||
keenetic_metrics: Dict[str, dict] = {}
|
keenetic_metrics: Dict[str, dict] = {}
|
||||||
_refresh_all_running = False
|
_refresh_all_running = False
|
||||||
|
_poll_lock = asyncio.Lock()
|
||||||
|
|
||||||
KEENETIC_FILE = DATA_DIR / "keenetic.json"
|
KEENETIC_FILE = DATA_DIR / "keenetic.json"
|
||||||
DEVICE_REFRESH_TIMEOUT = 75
|
DEVICE_REFRESH_TIMEOUT = 45
|
||||||
REFRESH_ALL_GAP_SEC = 2
|
REFRESH_ALL_GAP_SEC = 2
|
||||||
|
|
||||||
|
|
||||||
@@ -95,6 +96,11 @@ async def _refresh_device(dev: dict) -> dict:
|
|||||||
await client.close()
|
await client.close()
|
||||||
|
|
||||||
|
|
||||||
|
async def _refresh_device_locked(dev: dict) -> dict:
|
||||||
|
async with _poll_lock:
|
||||||
|
return await _refresh_device(dev)
|
||||||
|
|
||||||
|
|
||||||
@router.get("/list")
|
@router.get("/list")
|
||||||
async def keenetic_list(request: Request, user: str = Depends(require_auth)):
|
async def keenetic_list(request: Request, user: str = Depends(require_auth)):
|
||||||
devices = _load_keenetic()
|
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}
|
result = {"status": "ok", "web_url": web_url, "host": host}
|
||||||
if body.get("refresh", True):
|
if body.get("refresh", True):
|
||||||
try:
|
try:
|
||||||
metrics = await _refresh_device(devices[idx])
|
metrics = await _refresh_device_locked(devices[idx])
|
||||||
result["metrics"] = metrics
|
result["metrics"] = metrics
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
result["refresh_error"] = str(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"}
|
return {"status": "error", "detail": "router not found"}
|
||||||
|
|
||||||
try:
|
try:
|
||||||
metrics = await _refresh_device(dev)
|
metrics = await _refresh_device_locked(dev)
|
||||||
if not metrics["online"] and metrics.get("error"):
|
if not metrics["online"] and metrics.get("error"):
|
||||||
return {"status": "error", "detail": metrics["error"], "metrics": metrics}
|
return {"status": "error", "detail": metrics["error"], "metrics": metrics}
|
||||||
return {"status": "ok", "metrics": metrics}
|
return {"status": "ok", "metrics": metrics}
|
||||||
@@ -302,18 +308,19 @@ async def _refresh_all_devices() -> list:
|
|||||||
devices = _load_keenetic()
|
devices = _load_keenetic()
|
||||||
results = []
|
results = []
|
||||||
try:
|
try:
|
||||||
for i, dev in enumerate(devices):
|
async with _poll_lock:
|
||||||
try:
|
for i, dev in enumerate(devices):
|
||||||
metrics = await _refresh_device(dev)
|
try:
|
||||||
results.append({
|
metrics = await _refresh_device(dev)
|
||||||
"name": dev["name"],
|
results.append({
|
||||||
"online": metrics["online"],
|
"name": dev["name"],
|
||||||
"error": metrics.get("error", ""),
|
"online": metrics["online"],
|
||||||
})
|
"error": metrics.get("error", ""),
|
||||||
except Exception as e:
|
})
|
||||||
results.append({"name": dev["name"], "online": False, "error": str(e)})
|
except Exception as e:
|
||||||
if i + 1 < len(devices):
|
results.append({"name": dev["name"], "online": False, "error": str(e)})
|
||||||
await asyncio.sleep(REFRESH_ALL_GAP_SEC)
|
if i + 1 < len(devices):
|
||||||
|
await asyncio.sleep(REFRESH_ALL_GAP_SEC)
|
||||||
finally:
|
finally:
|
||||||
_refresh_all_running = False
|
_refresh_all_running = False
|
||||||
return results
|
return results
|
||||||
@@ -374,18 +381,19 @@ async def keenetic_reboot(name: str, request: Request, user: str = Depends(requi
|
|||||||
|
|
||||||
async def keenetic_monitor_loop():
|
async def keenetic_monitor_loop():
|
||||||
"""Background polling for all Keenetic routers."""
|
"""Background polling for all Keenetic routers."""
|
||||||
await asyncio.sleep(15)
|
await asyncio.sleep(5)
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
devices = _load_keenetic()
|
devices = _load_keenetic()
|
||||||
if devices:
|
if devices and not _refresh_all_running:
|
||||||
logger.info(f"Keenetic monitor: refreshing {len(devices)} routers")
|
logger.info(f"Keenetic monitor: refreshing {len(devices)} routers")
|
||||||
for dev in devices:
|
async with _poll_lock:
|
||||||
try:
|
for dev in devices:
|
||||||
await _refresh_device(dev)
|
try:
|
||||||
except Exception as e:
|
await _refresh_device(dev)
|
||||||
logger.error(f"Keenetic refresh {dev['name']}: {e}")
|
except Exception as e:
|
||||||
await asyncio.sleep(REFRESH_ALL_GAP_SEC)
|
logger.error(f"Keenetic refresh {dev['name']}: {e}")
|
||||||
|
await asyncio.sleep(REFRESH_ALL_GAP_SEC)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Keenetic monitor loop error: {e}")
|
logger.error(f"Keenetic monitor loop error: {e}")
|
||||||
|
|
||||||
|
|||||||
@@ -13,7 +13,8 @@ import aiohttp
|
|||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
KEENDNS_MARKERS = (".pro", ".club", ".link", "netcraze", "keenetic")
|
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}$")
|
_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)
|
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."""
|
"""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):
|
if is_public_ip_host(host):
|
||||||
return aiohttp.ClientTimeout(total=20, connect=8, sock_read=12)
|
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:
|
def build_api_base_url(host: str, web_url: str = "") -> str:
|
||||||
@@ -87,6 +111,7 @@ def _make_connector() -> aiohttp.TCPConnector:
|
|||||||
family=socket.AF_INET,
|
family=socket.AF_INET,
|
||||||
force_close=True,
|
force_close=True,
|
||||||
enable_cleanup_closed=True,
|
enable_cleanup_closed=True,
|
||||||
|
ttl_dns_cache=30,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -108,22 +133,47 @@ class KeeneticClient:
|
|||||||
self.base_url = base
|
self.base_url = base
|
||||||
parsed = urlparse(base)
|
parsed = urlparse(base)
|
||||||
self._host_key = parsed.netloc or host
|
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.login = login
|
||||||
self.password = password
|
self.password = password
|
||||||
self._session: Optional[aiohttp.ClientSession] = None
|
self._session: Optional[aiohttp.ClientSession] = None
|
||||||
self._authenticated = False
|
self._authenticated = False
|
||||||
self.last_error = ""
|
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:
|
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)
|
jar = aiohttp.CookieJar(unsafe=True)
|
||||||
self._session = aiohttp.ClientSession(
|
self._session = aiohttp.ClientSession(
|
||||||
timeout=timeout,
|
timeout=timeout,
|
||||||
cookie_jar=jar,
|
cookie_jar=jar,
|
||||||
connector=_make_connector(),
|
connector=_make_connector(),
|
||||||
version=aiohttp.HttpVersion11,
|
version=aiohttp.HttpVersion11,
|
||||||
headers={"User-Agent": "VPS-Monitoring/1.0"},
|
headers=self._request_headers(),
|
||||||
)
|
)
|
||||||
return self._session
|
return self._session
|
||||||
|
|
||||||
@@ -142,86 +192,103 @@ class KeeneticClient:
|
|||||||
return "Connection timeout"
|
return "Connection timeout"
|
||||||
|
|
||||||
async def _authenticate_once(self) -> bool:
|
async def _authenticate_once(self) -> bool:
|
||||||
"""Single auth attempt."""
|
"""Single auth attempt; tries each KeenDNS A record before giving up."""
|
||||||
self.last_error = ""
|
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 = (
|
retryable = (
|
||||||
aiohttp.ServerTimeoutError,
|
aiohttp.ServerTimeoutError,
|
||||||
aiohttp.ClientOSError,
|
aiohttp.ClientOSError,
|
||||||
asyncio.TimeoutError,
|
asyncio.TimeoutError,
|
||||||
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:
|
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:
|
except aiohttp.ClientConnectorError as e:
|
||||||
|
last_err = e
|
||||||
err = str(e).lower()
|
err = str(e).lower()
|
||||||
if "name or service not known" in err or "nodename nor servname" in err:
|
if "name or service not known" in err or "nodename nor servname" in err:
|
||||||
self.last_error = "DNS не резолвится с VPS"
|
self.last_error = "DNS не резолвится с VPS"
|
||||||
else:
|
return False
|
||||||
self.last_error = "Cannot connect to router"
|
continue
|
||||||
logger.error(f"Keenetic auth error: {type(e).__name__}: {e}")
|
|
||||||
return False
|
|
||||||
except retryable as e:
|
except retryable as e:
|
||||||
self.last_error = self._timeout_error_message()
|
last_err = e
|
||||||
logger.warning(
|
logger.debug(
|
||||||
f"Keenetic auth timeout ({attempt + 1}/{AUTH_RETRIES}) "
|
f"Keenetic auth try failed @ {auth_url}: {type(e).__name__}"
|
||||||
f"@ {self.base_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()
|
await self._reset_session()
|
||||||
if attempt + 1 < AUTH_RETRIES:
|
if await self._authenticate_once():
|
||||||
await asyncio.sleep(2.0 * (attempt + 1))
|
return True
|
||||||
continue
|
if self.last_error == "DNS не резолвится с VPS":
|
||||||
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
|
||||||
|
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
|
return False
|
||||||
|
|
||||||
async def rci_show(self, command: str, params: Optional[dict] = None) -> Optional[dict]:
|
async def rci_show(self, command: str, params: Optional[dict] = None) -> Optional[dict]:
|
||||||
@@ -232,7 +299,7 @@ class KeeneticClient:
|
|||||||
|
|
||||||
session = await self._get_session()
|
session = await self._get_session()
|
||||||
path = command.replace(" ", "/")
|
path = command.replace(" ", "/")
|
||||||
url = f"{self.base_url}/rci/show/{path}"
|
url = f"{self._effective_base()}/rci/show/{path}"
|
||||||
|
|
||||||
try:
|
try:
|
||||||
async with session.get(url, params=params, ssl=False) as resp:
|
async with session.get(url, params=params, ssl=False) as resp:
|
||||||
@@ -262,7 +329,7 @@ class KeeneticClient:
|
|||||||
asyncio.TimeoutError,
|
asyncio.TimeoutError,
|
||||||
TimeoutError,
|
TimeoutError,
|
||||||
)
|
)
|
||||||
url = f"{self.base_url}/rci/"
|
url = f"{self._effective_base()}/rci/"
|
||||||
for attempt in range(AUTH_RETRIES):
|
for attempt in range(AUTH_RETRIES):
|
||||||
try:
|
try:
|
||||||
session = await self._get_session()
|
session = await self._get_session()
|
||||||
|
|||||||
Reference in New Issue
Block a user