OSV states Debian ranges in *source* package versions. A source package that ships several binaries gives each its own upstream version, and the two are unrelated numbers: `libldb2` is `2:2.11.0+samba4.22.11+dfsg-0+deb13u1` while its source, samba, is `2:4.22.11+dfsg-…`. A consumer that resolves the coordinate on `source_package` — which is what this driver's `source:Package` is for — and then compares `version` is comparing ldb's version against samba's range. dpkg reads `2.11.0` as older than the `2:4.17.4+dfsg-1` that fixed CVE-2022-44640, so a Debian 13 host running samba 4.22.11, five releases past the fix, was reported vulnerable on four packages at once. This driver cannot fix that comparison. It is the only place that can supply the number to make it with. Empty when dpkg considers it equal to `Version`, and empty on a dpkg that does not know the field, so it falls back to `Version` — which is the previous behaviour and correct everywhere except the shape above. `maxsplit` goes from 4 to 5 with the extra field. Summary stays last, so it keeps whatever it contains. **The fixture was carrying four fields against a format string asking for five.** `source_package` had been silently receiving the description, and no assertion looked at it. It now carries what dpkg-query actually returns, including a package whose source version is a different number from its own. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2350 lines
96 KiB
Python
2350 lines
96 KiB
Python
# -*- coding: utf-8 -*-
|
||
# Licensed under the Apache License, Version 2.0
|
||
|
||
"""NAPALM driver for generic Linux systems.
|
||
|
||
Connects via SSH using netmiko (device_type ``linux``) and supports
|
||
auto-detection of the installed package manager:
|
||
|
||
* apt — Debian, Ubuntu, Raspberry Pi OS, …
|
||
* dnf — RHEL 8+, Rocky Linux, AlmaLinux, Fedora
|
||
* yum — RHEL 7, CentOS 7
|
||
* apk — Alpine Linux
|
||
* pacman — Arch Linux, Manjaro
|
||
|
||
A specific package manager can be forced with
|
||
``optional_args={"pkg_manager": "apt"}``.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import re
|
||
import socket
|
||
from shlex import quote as _shlex_quote
|
||
from typing import Any
|
||
|
||
from netmiko import ConnectHandler
|
||
from netmiko.exceptions import (
|
||
NetmikoAuthenticationException,
|
||
NetmikoTimeoutException,
|
||
)
|
||
from napalm.base.exceptions import ConnectionException, ConnectionClosedException
|
||
from napalm.base.netmiko_helpers import netmiko_args
|
||
from napalm_device_types import FingerprintRule, OSDriver
|
||
from napalm_device_types.models import (
|
||
ApplyUpdatesResultDict,
|
||
CronJobDict,
|
||
DeviceActionResultDict,
|
||
DockerInfoDict,
|
||
PackageDict,
|
||
ProcessDict,
|
||
ServiceDict,
|
||
SNMPConfigDict,
|
||
UpdateDict,
|
||
UserDict,
|
||
)
|
||
|
||
logger = logging.getLogger("napalm_linux")
|
||
|
||
# Package managers in detection order
|
||
_PKG_MANAGERS = ["apt", "dnf", "yum", "apk", "pacman"]
|
||
|
||
# DMI field values that carry no useful information (OEM defaults, blanks)
|
||
_BAD_DMI: frozenset[str] = frozenset({
|
||
"", "none", "n/a", "not specified", "not applicable",
|
||
"to be filled by o.e.m.", "default string", "unknown",
|
||
"no asset tag", "not present",
|
||
})
|
||
|
||
# systemd-detect-virt output → human-readable vendor name
|
||
_VIRT_VENDOR_MAP: dict[str, str] = {
|
||
"kvm": "KVM",
|
||
"qemu": "KVM",
|
||
"vmware": "VMware ESXi",
|
||
"microsoft": "Microsoft Hyper-V",
|
||
"xen": "Xen",
|
||
"virtualbox": "Oracle VirtualBox",
|
||
"parallels": "Parallels",
|
||
"docker": "Docker",
|
||
"podman": "Podman",
|
||
"lxc": "LXC",
|
||
"lxc-libvirt": "LXC",
|
||
"systemd-nspawn": "systemd-nspawn",
|
||
}
|
||
|
||
# Container technologies reported by systemd-detect-virt
|
||
_CONTAINER_VIRT: frozenset[str] = frozenset({
|
||
"docker", "podman", "lxc", "lxc-libvirt", "systemd-nspawn",
|
||
})
|
||
|
||
# DMI sys_vendor strings that indicate a VM when detect-virt is unavailable
|
||
_VM_DMI_VENDORS: frozenset[str] = frozenset({
|
||
"qemu", "vmware, inc.", "microsoft corporation",
|
||
"innotek gmbh", "xen", "bochs",
|
||
"parallels software international inc.",
|
||
})
|
||
|
||
# Known ARM board model prefixes → canonical vendor name
|
||
_ARM_VENDOR_PREFIXES: list[tuple[str, str]] = [
|
||
("Raspberry Pi", "Raspberry Pi Foundation"),
|
||
("NVIDIA Jetson", "NVIDIA"),
|
||
("ODROID", "Hardkernel"),
|
||
("Hardkernel", "Hardkernel"),
|
||
("Rock Pi", "Radxa"),
|
||
("ROCK Pi", "Radxa"),
|
||
("Radxa", "Radxa"),
|
||
("Orange Pi", "Xunlong Software"),
|
||
("Banana Pi", "SinoVoip"),
|
||
("NanoPi", "FriendlyElec"),
|
||
("PINE64", "Pine64"),
|
||
("BeagleBone", "BeagleBoard.org"),
|
||
]
|
||
|
||
|
||
def _arm_vendor_from_model(model: str) -> str:
|
||
"""Extract vendor from an ARM device tree / cpuinfo model string."""
|
||
for prefix, vendor in _ARM_VENDOR_PREFIXES:
|
||
if model.startswith(prefix):
|
||
return vendor
|
||
# Generic fallback: words before the first numeric token
|
||
brand = []
|
||
for word in model.split():
|
||
if word[0].isdigit():
|
||
break
|
||
brand.append(word)
|
||
return " ".join(brand)
|
||
|
||
|
||
#: A bare image ID (``d626d04934cd``), not a registry reference. ``docker ps``
|
||
#: falls back to this whenever the tag a container was created from has since
|
||
#: been moved to a newer image — i.e. exactly after a pull without a recreate.
|
||
_IMAGE_ID_RE = re.compile(r"^(sha256:)?[0-9a-f]{12,64}$")
|
||
|
||
|
||
def _looks_like_image_id(ref: str) -> bool:
|
||
"""True if *ref* is an image ID rather than something a registry can resolve."""
|
||
return bool(_IMAGE_ID_RE.match(ref.strip()))
|
||
|
||
|
||
def _short_image_id(raw: str) -> str:
|
||
"""Normalise ``sha256:<64hex>`` and ``<12hex>`` to a comparable 12-char form."""
|
||
return raw.strip().removeprefix("sha256:")[:12]
|
||
|
||
|
||
class LinuxDriver(OSDriver):
|
||
"""NAPALM driver for generic Linux systems.
|
||
|
||
Connects via SSH (netmiko ``linux`` device type) and auto-detects the
|
||
package manager unless overridden by ``optional_args["pkg_manager"]``.
|
||
"""
|
||
|
||
TYPE_LABEL = "Linux"
|
||
VENDOR = "Linux"
|
||
DRIVER_NAME = "linux"
|
||
# A general-purpose host runs through a full init sequence; NAS derivatives
|
||
# (OpenMediaVault, QNAP) inherit this and are, if anything, slower.
|
||
REBOOT_SETTLE_SECONDS = 90
|
||
SNMP_FINGERPRINT = [
|
||
FingerprintRule("linux", weight=5.0),
|
||
]
|
||
SSH_FINGERPRINT = [
|
||
FingerprintRule("ubuntu", weight=4.0),
|
||
FingerprintRule("debian", weight=4.0),
|
||
FingerprintRule("alpine", weight=4.0),
|
||
FingerprintRule("openssh", weight=1.0),
|
||
]
|
||
NETMIKO_DEVICE_TYPE = "linux"
|
||
|
||
def __init__(
|
||
self,
|
||
hostname: str,
|
||
username: str,
|
||
password: str,
|
||
timeout: int = 60,
|
||
optional_args: dict | None = None,
|
||
) -> None:
|
||
self.hostname = hostname
|
||
self.username = username
|
||
self.password = password
|
||
self.timeout = timeout
|
||
|
||
if optional_args is None:
|
||
optional_args = {}
|
||
|
||
self.port: int = optional_args.get("port", 22)
|
||
self._forced_pkg_manager: str | None = optional_args.get("pkg_manager")
|
||
self._secret: str = optional_args.get("secret", password)
|
||
# Optional sudo password for privilege escalation (e.g. apt-get update)
|
||
self._sudo_password: str | None = optional_args.get("sudo_password")
|
||
# Expected apt proxy URL — checked as a device warning on apt systems.
|
||
# Only set when apt_proxy_enabled=true in NetOrk settings; empty string disables the check.
|
||
self._apt_proxy_url: str = optional_args.get("apt_proxy_url", "")
|
||
|
||
if optional_args.get("debugging"):
|
||
logger.setLevel(logging.DEBUG)
|
||
|
||
self.netmiko_optional_args = netmiko_args(optional_args)
|
||
# port is passed explicitly in open() — remove it from netmiko_optional_args
|
||
# to avoid "multiple values for keyword argument 'port'"
|
||
self.netmiko_optional_args.pop("port", None)
|
||
|
||
# Runtime state
|
||
self._device: ConnectHandler | None = None
|
||
self._pkg_manager: str | None = None # set after open()
|
||
|
||
# ------------------------------------------------------------------
|
||
# Connection management
|
||
# ------------------------------------------------------------------
|
||
|
||
def open(self) -> None:
|
||
"""Open the SSH connection and detect the package manager."""
|
||
try:
|
||
self._device = ConnectHandler(
|
||
device_type=self.NETMIKO_DEVICE_TYPE,
|
||
host=self.hostname,
|
||
username=self.username,
|
||
password=self.password,
|
||
port=self.port,
|
||
secret=self._secret,
|
||
timeout=self.timeout,
|
||
**self.netmiko_optional_args,
|
||
)
|
||
except NetmikoAuthenticationException as exc:
|
||
raise ConnectionException(str(exc)) from exc
|
||
except NetmikoTimeoutException as exc:
|
||
raise ConnectionException(str(exc)) from exc
|
||
|
||
# Prevent PTY from wrapping long output lines (e.g. docker JSON).
|
||
try:
|
||
self._device.send_command("stty cols 10000 2>/dev/null || true", expect_string=r"[#$>]\s*$")
|
||
except Exception:
|
||
pass
|
||
|
||
self._pkg_manager = self._forced_pkg_manager or self._detect_pkg_manager()
|
||
logger.debug("Connected to %s, pkg_manager=%s", self.hostname, self._pkg_manager)
|
||
|
||
def close(self) -> None:
|
||
"""Close the SSH connection."""
|
||
if self._device:
|
||
try:
|
||
self._device.disconnect()
|
||
except Exception:
|
||
pass
|
||
self._device = None
|
||
self._pkg_manager = None
|
||
|
||
def is_alive(self) -> dict[str, bool]:
|
||
if self._device:
|
||
try:
|
||
return {"is_alive": self._device.remote_conn.transport.is_active()}
|
||
except (AttributeError, socket.error, EOFError):
|
||
return {"is_alive": False}
|
||
return {"is_alive": False}
|
||
|
||
# ------------------------------------------------------------------
|
||
# Internal helpers
|
||
# ------------------------------------------------------------------
|
||
|
||
def _send(self, command: str, read_timeout: float = 100) -> str:
|
||
"""Send a command and return stripped output."""
|
||
if not self._device:
|
||
raise ConnectionClosedException("Not connected")
|
||
return self._device.send_command(
|
||
command,
|
||
read_timeout=read_timeout,
|
||
cmd_verify=False,
|
||
expect_string=r'[#$\>]\s*$',
|
||
).strip()
|
||
|
||
def _sudo(self, command: str, read_timeout: float = 100) -> str:
|
||
"""Run *command* via sudo, feeding the password via stdin (-S).
|
||
|
||
Falls back to plain execution when no sudo password is configured.
|
||
"""
|
||
if self._sudo_password:
|
||
wrapped = f'echo {_shlex_quote(self._sudo_password)} | sudo -S -p "" {command}'
|
||
return self._send(wrapped, read_timeout=read_timeout)
|
||
return self._send(f'sudo {command}', read_timeout=read_timeout)
|
||
|
||
def _detect_pkg_manager(self) -> str | None:
|
||
"""Return the first package manager binary found on PATH."""
|
||
for pm in _PKG_MANAGERS:
|
||
result = self._send(f"command -v {pm} 2>/dev/null")
|
||
if result:
|
||
return pm
|
||
return None
|
||
|
||
# ------------------------------------------------------------------
|
||
# Standard NAPALM – read-only
|
||
# ------------------------------------------------------------------
|
||
|
||
def _collect_platform_info(self) -> dict[str, Any]:
|
||
"""Collect hardware/virtualisation info in a single SSH round-trip.
|
||
|
||
Returns a dict with keys:
|
||
- vendor (str) — hardware vendor or hypervisor name; "" if unknown
|
||
- model (str) — product model or "Virtual Machine"/"Container"; "" if unknown
|
||
- serial (str) — product serial, or VM UUID as fallback; "" if unknown
|
||
- is_vm (bool) — True for VMs and containers
|
||
"""
|
||
dmi_cmd = (
|
||
"v=$(cat /sys/class/dmi/id/sys_vendor 2>/dev/null); "
|
||
"n=$(cat /sys/class/dmi/id/product_name 2>/dev/null); "
|
||
"r=$(cat /sys/class/dmi/id/product_version 2>/dev/null); "
|
||
"s=$(cat /sys/class/dmi/id/product_serial 2>/dev/null); "
|
||
"u=$(cat /sys/class/dmi/id/product_uuid 2>/dev/null); "
|
||
"d=$(systemd-detect-virt 2>/dev/null); d=${d:-none}; "
|
||
"dt=$(tr -d '\\0' </sys/firmware/devicetree/base/model 2>/dev/null); "
|
||
"cs=$(grep '^Serial' /proc/cpuinfo 2>/dev/null | head -1 | cut -d: -f2 | xargs 2>/dev/null); "
|
||
"cm=$(grep '^Model' /proc/cpuinfo 2>/dev/null | head -1 | cut -d: -f2 | xargs 2>/dev/null); "
|
||
# DMIBEGIN sentinel: _send() strips leading blank lines (ARM has no DMI
|
||
# files, so fields 0-4 are empty). The sentinel anchors the output so
|
||
# splitlines()[start+N] always maps to the correct field index.
|
||
"printf 'DMIBEGIN\\n%s\\n%s\\n%s\\n%s\\n%s\\n%s\\n%s\\n%s\\n%s\\n' "
|
||
"\"$v\" \"$n\" \"$r\" \"$s\" \"$u\" \"$d\" \"$dt\" \"$cs\" \"$cm\""
|
||
)
|
||
try:
|
||
raw_lines = self._send(dmi_cmd).splitlines()
|
||
try:
|
||
start = raw_lines.index("DMIBEGIN") + 1
|
||
except ValueError:
|
||
start = 0
|
||
lines = raw_lines[start:]
|
||
except Exception:
|
||
return {"vendor": "", "model": "", "serial": "", "is_vm": False}
|
||
|
||
def _clean(idx: int) -> str:
|
||
val = lines[idx].strip() if idx < len(lines) else ""
|
||
return "" if val.lower() in _BAD_DMI else val
|
||
|
||
sys_vendor = _clean(0)
|
||
product_name = _clean(1)
|
||
product_ver = _clean(2)
|
||
product_ser = _clean(3)
|
||
product_uuid = _clean(4)
|
||
detect_virt = lines[5].strip().lower() if len(lines) > 5 else "none"
|
||
dt_model = _clean(6)
|
||
cpuinfo_ser = _clean(7)
|
||
cpuinfo_mdl = _clean(8)
|
||
|
||
is_container = detect_virt in _CONTAINER_VIRT
|
||
is_vm = (
|
||
detect_virt not in ("none", "")
|
||
or sys_vendor.lower() in _VM_DMI_VENDORS
|
||
)
|
||
|
||
if is_container:
|
||
return {
|
||
"vendor": _VIRT_VENDOR_MAP.get(detect_virt, sys_vendor or "Container"),
|
||
"model": "Container",
|
||
"serial": product_uuid,
|
||
"is_vm": True,
|
||
}
|
||
|
||
if is_vm:
|
||
vendor = _VIRT_VENDOR_MAP.get(detect_virt, "")
|
||
if not vendor:
|
||
sv = sys_vendor.lower()
|
||
if "vmware" in sv:
|
||
vendor = "VMware ESXi"
|
||
elif "microsoft" in sv:
|
||
vendor = "Microsoft Hyper-V"
|
||
elif "qemu" in sv or "kvm" in sv:
|
||
vendor = "KVM"
|
||
elif "xen" in sv:
|
||
vendor = "Xen"
|
||
elif "innotek" in sv or "virtualbox" in sv:
|
||
vendor = "Oracle VirtualBox"
|
||
else:
|
||
vendor = sys_vendor
|
||
return {
|
||
"vendor": vendor,
|
||
"model": "Virtual Machine",
|
||
"serial": product_ser or product_uuid,
|
||
"is_vm": True,
|
||
}
|
||
|
||
# Bare-metal: prefer product_version when it reads like a marketing name
|
||
if sys_vendor or product_name:
|
||
pv_usable = product_ver and product_ver != product_name and " " in product_ver
|
||
return {
|
||
"vendor": sys_vendor,
|
||
"model": product_ver if pv_usable else product_name,
|
||
"serial": product_ser,
|
||
"is_vm": False,
|
||
}
|
||
|
||
# ARM/embedded fallback: no DMI, try device tree and /proc/cpuinfo
|
||
arm_model = dt_model or cpuinfo_mdl
|
||
if arm_model:
|
||
return {
|
||
"vendor": _arm_vendor_from_model(arm_model),
|
||
"model": arm_model,
|
||
"serial": cpuinfo_ser,
|
||
"is_vm": False,
|
||
}
|
||
|
||
return {"vendor": "", "model": "", "serial": "", "is_vm": False}
|
||
|
||
def get_facts(self) -> dict[str, Any]:
|
||
"""Return basic system facts."""
|
||
hostname = self._send("hostname -s 2>/dev/null || hostname")
|
||
fqdn = self._send("hostname -f 2>/dev/null || hostname")
|
||
os_version = self._send(
|
||
"cat /etc/os-release 2>/dev/null | grep '^PRETTY_NAME' | cut -d= -f2 | tr -d '\"'"
|
||
) or self._send("uname -r")
|
||
uptime_secs = self._parse_uptime()
|
||
platform = self._collect_platform_info()
|
||
|
||
iface_out = self._send("ip -o link show | awk -F': ' '{print $2}' | cut -d@ -f1")
|
||
interface_list = [i.strip() for i in iface_out.splitlines() if i.strip() and i.strip() != "lo"]
|
||
|
||
# Currently-booted kernel release, distinct from an installed-but-not-yet-
|
||
# booted newer kernel (used for kernel CVE relevance).
|
||
running_kernel = self._send("uname -r").strip()
|
||
|
||
return {
|
||
"hostname": hostname,
|
||
"fqdn": fqdn,
|
||
"vendor": platform["vendor"] or self.VENDOR,
|
||
"model": platform["model"],
|
||
"serial_number": platform["serial"],
|
||
"os_version": os_version,
|
||
"uptime": uptime_secs,
|
||
"interface_list": interface_list,
|
||
"running_kernel": running_kernel,
|
||
}
|
||
|
||
def _parse_uptime(self) -> int:
|
||
"""Return uptime in seconds from ``/proc/uptime``."""
|
||
raw = self._send("cat /proc/uptime 2>/dev/null")
|
||
try:
|
||
return int(float(raw.split()[0]))
|
||
except (IndexError, ValueError):
|
||
return 0
|
||
|
||
def get_lldp_neighbors(self) -> Dict[str, List[dict[str, Any]]]:
|
||
"""Return LLDP neighbors if lldpd is installed and currently running.
|
||
|
||
Uses ``lldpctl -f keyvalue``. Returns an empty dict when lldpd is
|
||
absent or stopped — does NOT attempt to start the daemon.
|
||
"""
|
||
# Check lldpctl is available
|
||
if not self._send("command -v lldpctl 2>/dev/null").strip():
|
||
return {}
|
||
|
||
# Check lldpd is active (systemd or fallback to pgrep)
|
||
running = self._send(
|
||
"systemctl is-active lldpd 2>/dev/null || "
|
||
"service lldpd status 2>/dev/null | grep -q running && echo active || "
|
||
"pgrep -x lldpd >/dev/null 2>&1 && echo active || true"
|
||
).strip()
|
||
if "active" not in running:
|
||
return {}
|
||
|
||
output = self._send("lldpctl -f keyvalue 2>/dev/null || true")
|
||
neighbors: Dict[str, List[dict[str, Any]]] = {}
|
||
entries: Dict[str, Dict[str, str]] = {}
|
||
|
||
for line in output.splitlines():
|
||
line = line.strip()
|
||
if "=" not in line:
|
||
continue
|
||
key, _, value = line.partition("=")
|
||
parts = key.split(".")
|
||
if len(parts) < 3 or parts[0] != "lldp":
|
||
continue
|
||
iface = parts[1]
|
||
subkey = ".".join(parts[2:])
|
||
entries.setdefault(iface, {})[subkey] = value
|
||
|
||
for iface, data in entries.items():
|
||
entry = {
|
||
"hostname": data.get("chassis.name", ""),
|
||
"port": data.get("port.ifname", data.get("port.id.value", "")),
|
||
}
|
||
if data.get("chassis.id.subtype") == "mac":
|
||
mac = data.get("chassis.id.value", "")
|
||
if re.match(r"^([0-9A-Fa-f]{2}:){5}[0-9A-Fa-f]{2}$", mac):
|
||
entry["mac"] = mac.lower()
|
||
neighbors.setdefault(iface, []).append(entry)
|
||
|
||
return neighbors
|
||
|
||
def get_interfaces(self) -> dict[str, Any]:
|
||
"""Return interface operational data."""
|
||
interfaces: dict[str, Any] = {}
|
||
|
||
# ip -o link show: one line per interface
|
||
link_out = self._send("ip -o link show")
|
||
for line in link_out.splitlines():
|
||
# 2: eth0: <BROADCAST,MULTICAST,UP,LOWER_UP> mtu 1500 ... state UP
|
||
m = re.match(r"^\d+:\s+(\S+?)(?:@\S+)?:\s+<([^>]*)>.*mtu\s+(\d+).*state\s+(\S+)", line)
|
||
if not m:
|
||
continue
|
||
name, flags, mtu, state = m.group(1), m.group(2), int(m.group(3)), m.group(4)
|
||
mac_m = re.search(r"link/ether\s+([\da-f:]+)", line)
|
||
mac = mac_m.group(1) if mac_m else ""
|
||
is_up = "UP" in flags.split(",") or state == "UP"
|
||
interfaces[name] = {
|
||
"is_up": is_up,
|
||
"is_enabled": "UP" in flags.split(","),
|
||
"description": "",
|
||
"last_flapped": -1.0,
|
||
"speed": -1.0,
|
||
"mtu": mtu,
|
||
"mac_address": mac,
|
||
}
|
||
|
||
return interfaces
|
||
|
||
def get_interfaces_ip(self) -> dict[str, Any]:
|
||
"""Return IP addresses per interface."""
|
||
result: dict[str, Any] = {}
|
||
|
||
addr_out = self._send("ip -o addr show")
|
||
for line in addr_out.splitlines():
|
||
# 2: eth0 inet 192.168.1.10/24 brd ...
|
||
m = re.match(r"^\d+:\s+(\S+)\s+(inet6?)\s+([\da-f.:]+)/(\d+)", line)
|
||
if not m:
|
||
continue
|
||
iface, family, addr, prefix = m.group(1), m.group(2), m.group(3), int(m.group(4))
|
||
af = "ipv4" if family == "inet" else "ipv6"
|
||
result.setdefault(iface, {"ipv4": {}, "ipv6": {}})
|
||
result[iface][af][addr] = {"prefix_length": prefix}
|
||
|
||
return result
|
||
|
||
def get_networks(self) -> list[dict[str, Any]]:
|
||
"""Return IP networks derived from interface addresses.
|
||
|
||
Excludes loopback, link-local, /32 host-only addresses, and
|
||
container-internal interfaces (docker*, br-*, veth*, virbr*).
|
||
|
||
Each entry matches the OPNsense get_networks() schema::
|
||
|
||
{
|
||
"network": "10.7.224.0/24",
|
||
"interface": "ens7",
|
||
"gateway": "10.7.224.11",
|
||
"family": "ipv4",
|
||
"prefix_length": 24,
|
||
"vlan_id": None,
|
||
}
|
||
"""
|
||
import ipaddress
|
||
|
||
_SKIP_PREFIXES = ("lo", "docker", "br-", "veth", "virbr", "tun", "tap")
|
||
networks: list[dict[str, Any]] = []
|
||
|
||
for iface_name, af_data in self.get_interfaces_ip().items():
|
||
if any(iface_name.startswith(p) for p in _SKIP_PREFIXES):
|
||
continue
|
||
for family, addrs in af_data.items():
|
||
for addr, info in addrs.items():
|
||
prefix = info.get("prefix_length", 0)
|
||
try:
|
||
iface_obj = ipaddress.ip_interface(f"{addr}/{prefix}")
|
||
net = iface_obj.network
|
||
if net.is_loopback or net.is_link_local:
|
||
continue
|
||
# Skip host-only addresses (/32 IPv4, /128 IPv6)
|
||
if (net.version == 4 and net.prefixlen >= 32) or (
|
||
net.version == 6 and net.prefixlen >= 128
|
||
):
|
||
continue
|
||
networks.append({
|
||
"network": str(net),
|
||
"interface": iface_name,
|
||
"gateway": str(iface_obj.ip),
|
||
"family": family,
|
||
"prefix_length": net.prefixlen,
|
||
"vlan_id": None,
|
||
})
|
||
except ValueError:
|
||
pass
|
||
|
||
return networks
|
||
|
||
def get_route_to(
|
||
self,
|
||
destination: str = "",
|
||
protocol: str = "",
|
||
longer: bool = False,
|
||
) -> Dict[str, List[dict[str, Any]]]:
|
||
"""Return the routing table via ``ip route show``.
|
||
|
||
OSPF routes (from FRR/Quagga) are included via ``ip route show proto ospf``
|
||
if any are present. The result is keyed by network prefix.
|
||
"""
|
||
routes: Dict[str, List[dict[str, Any]]] = {}
|
||
|
||
proto_map = {
|
||
"kernel": "connected",
|
||
"dhcp": "dhcp",
|
||
"static": "static",
|
||
"ospf": "ospf",
|
||
"bgp": "bgp",
|
||
"bird": "bgp",
|
||
"ra": "connected",
|
||
"boot": "connected",
|
||
"zebra": "zebra",
|
||
}
|
||
|
||
def _make_entry(proto: str, nexthop: str, iface: str, metric: int, network: str) -> dict[str, Any]:
|
||
family = "ipv6" if (":" in network or (nexthop and ":" in nexthop)) else "ipv4"
|
||
return {
|
||
"protocol": proto,
|
||
"family": family,
|
||
"current_active": True,
|
||
"last_active": False,
|
||
"age": -1,
|
||
"next_hop": nexthop,
|
||
"outgoing_interface": iface,
|
||
"selected_next_hop": True,
|
||
"preference": metric,
|
||
"routing_table": "global",
|
||
"protocol_attributes": {},
|
||
}
|
||
|
||
def _add(network: str, proto: str, nexthop: str, iface: str, metric: int) -> None:
|
||
if destination and network != destination:
|
||
return
|
||
mapped = proto_map.get(proto, proto)
|
||
if protocol and mapped != protocol.lower():
|
||
return
|
||
routes.setdefault(network, []).append(
|
||
_make_entry(mapped, nexthop, iface, metric, network)
|
||
)
|
||
|
||
out = self._send("ip -4 route show && ip -6 route show")
|
||
for line in out.splitlines():
|
||
line = line.strip()
|
||
if not line or line.startswith("#"):
|
||
continue
|
||
dest_m = re.match(r"^(\S+)", line)
|
||
if not dest_m:
|
||
continue
|
||
raw_dest = dest_m.group(1)
|
||
network = "0.0.0.0/0" if raw_dest == "default" else ("::/0" if raw_dest == "default6" else raw_dest)
|
||
if "/" not in network:
|
||
network += "/32"
|
||
|
||
nexthop = ""
|
||
nh_m = re.search(r"via\s+(\S+)", line)
|
||
if nh_m:
|
||
nexthop = nh_m.group(1)
|
||
|
||
iface = ""
|
||
dev_m = re.search(r"dev\s+(\S+)", line)
|
||
if dev_m:
|
||
iface = dev_m.group(1)
|
||
|
||
proto = "kernel"
|
||
proto_m = re.search(r"proto\s+(\S+)", line)
|
||
if proto_m:
|
||
proto = proto_m.group(1)
|
||
|
||
metric = 0
|
||
metric_m = re.search(r"metric\s+(\d+)", line)
|
||
if metric_m:
|
||
metric = int(metric_m.group(1))
|
||
|
||
_add(network, proto, nexthop, iface, metric)
|
||
|
||
# FRR/Zebra enrichment via vtysh — properly attributes OSPF/BGP/RIP protocols.
|
||
# FRR installs routes into the kernel as "proto zebra"; vtysh gives the real source.
|
||
_frr_code: Dict[str, str] = {
|
||
"O": "ospf", "B": "bgp", "R": "rip", "I": "isis",
|
||
"S": "static", "K": "connected", "C": "connected",
|
||
}
|
||
if self._send("command -v vtysh 2>/dev/null").strip():
|
||
try:
|
||
vtysh_out = self._send(
|
||
"vtysh -c 'show ip route' 2>/dev/null; vtysh -c 'show ipv6 route' 2>/dev/null",
|
||
read_timeout=15,
|
||
)
|
||
for vline in vtysh_out.splitlines():
|
||
# "O>* 10.10.0.0/24 [110/20] via 10.255.255.2, wg0, ..."
|
||
vm = re.match(
|
||
r"^([OBSCRIKEF])[>*\s]{0,3}([\d.:a-fA-F/]+)\s+\[(\d+)/(\d+)\]"
|
||
r"(?:\s+via\s+([\d.:a-fA-F]+),\s*(\S+?)(?:,|$))?",
|
||
vline.strip(),
|
||
)
|
||
if not vm:
|
||
continue
|
||
code = vm.group(1)
|
||
prefix = vm.group(2)
|
||
metric = int(vm.group(4))
|
||
nexthop = vm.group(5) or ""
|
||
iface = (vm.group(6) or "").rstrip(",")
|
||
if "/" not in prefix:
|
||
prefix += "/32"
|
||
mapped_proto = _frr_code.get(code, code.lower())
|
||
if prefix in routes:
|
||
for entry in routes[prefix]:
|
||
entry["protocol"] = mapped_proto
|
||
else:
|
||
if not (destination and prefix != destination):
|
||
routes.setdefault(prefix, []).append(
|
||
_make_entry(mapped_proto, nexthop, iface, metric, prefix)
|
||
)
|
||
except Exception:
|
||
pass
|
||
|
||
return routes
|
||
|
||
def get_arp_table(self, vrf: str = "") -> List[dict[str, Any]]:
|
||
"""Return the ARP/neighbour table."""
|
||
entries = []
|
||
neigh_out = self._send("ip -4 neigh show")
|
||
for line in neigh_out.splitlines():
|
||
# 192.168.1.1 dev eth0 lladdr aa:bb:cc:dd:ee:ff REACHABLE
|
||
m = re.match(
|
||
r"^([\d.]+)\s+dev\s+(\S+)\s+lladdr\s+([\da-f:]+)\s+(\S+)", line
|
||
)
|
||
if not m:
|
||
continue
|
||
entries.append({
|
||
"interface": m.group(2),
|
||
"mac": m.group(3),
|
||
"ip": m.group(1),
|
||
"age": 0.0,
|
||
})
|
||
return entries
|
||
|
||
def get_config(
|
||
self, retrieve: str = "all", full: bool = False, sanitized: bool = False
|
||
) -> dict[str, Any]:
|
||
"""Return minimal config representation (network interfaces only).
|
||
|
||
``valid_lft`` / ``preferred_lft`` fields from DHCP leases are stripped
|
||
so the output is stable across polls and does not produce false-positive
|
||
config-change events in the config-backup feature. ``veth*`` interfaces
|
||
are stripped entirely — Docker creates/destroys them with a fresh index
|
||
and random name on every container restart, which would otherwise flag
|
||
a config change on nearly every poll of a Docker host.
|
||
"""
|
||
import re
|
||
|
||
running = self._send("ip addr show && ip route show")
|
||
# Strip volatile DHCP lease timer fields — they decrement every poll
|
||
running = re.sub(r"\s+valid_lft\s+\S+\s+preferred_lft\s+\S+", "", running)
|
||
# Strip veth interface blocks (line + indented sub-lines) — ephemeral
|
||
# Docker container network endpoints, not intentional host config
|
||
running = re.sub(r"^\d+: veth\S*:.*\n(?:[ \t]+.*\n?)*", "", running, flags=re.MULTILINE)
|
||
return {"running": running, "startup": "", "candidate": ""}
|
||
|
||
# ------------------------------------------------------------------
|
||
# Config management – not applicable for generic Linux
|
||
# ------------------------------------------------------------------
|
||
|
||
def load_merge_candidate(self, filename: str = None, config: str = None) -> None: # type: ignore[override]
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
def load_replace_candidate(self, filename: str = None, config: str = None) -> None: # type: ignore[override]
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
def compare_config(self) -> str:
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
def commit_config(self, message: str = "") -> None:
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
def discard_config(self) -> None:
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
def rollback(self) -> None:
|
||
raise NotImplementedError("Config management is not supported for Linux hosts")
|
||
|
||
# ------------------------------------------------------------------
|
||
# Optional NAPALM methods
|
||
# ------------------------------------------------------------------
|
||
|
||
def ping(
|
||
self,
|
||
destination: str,
|
||
source: str = "",
|
||
ttl: int = 255,
|
||
timeout: int = 2,
|
||
size: int = 100,
|
||
count: int = 5,
|
||
vrf: str = "",
|
||
) -> dict[str, Any]:
|
||
"""Execute ping from the remote host."""
|
||
src_opt = f"-I {source}" if source else ""
|
||
cmd = f"ping -c {count} -W {timeout} -s {size} -t {ttl} {src_opt} {destination} 2>&1"
|
||
output = self._send(cmd)
|
||
|
||
# Parse summary line: "5 packets transmitted, 5 received, 0% packet loss"
|
||
m = re.search(
|
||
r"(\d+) packets transmitted,\s*(\d+) received,\s*([\d.]+)% packet loss",
|
||
output,
|
||
)
|
||
if not m:
|
||
return {"error": output}
|
||
|
||
sent, received = int(m.group(1)), int(m.group(2))
|
||
|
||
# Parse rtt line: "rtt min/avg/max/mdev = 0.123/0.456/0.789/0.100 ms"
|
||
rtt_m = re.search(
|
||
r"rtt .* = ([\d.]+)/([\d.]+)/([\d.]+)/([\d.]+) ms", output
|
||
)
|
||
|
||
results = []
|
||
for line in output.splitlines():
|
||
icmp_m = re.search(
|
||
r"bytes from ([\d.]+).*icmp_seq=\d+ ttl=(\d+) time=([\d.]+) ms", line
|
||
)
|
||
if icmp_m:
|
||
results.append({
|
||
"ip_address": icmp_m.group(1),
|
||
"rtt": float(icmp_m.group(3)),
|
||
})
|
||
|
||
return {
|
||
"success": {
|
||
"probes_sent": sent,
|
||
"packet_loss": sent - received,
|
||
"rtt_min": float(rtt_m.group(1)) if rtt_m else 0.0,
|
||
"rtt_avg": float(rtt_m.group(2)) if rtt_m else 0.0,
|
||
"rtt_max": float(rtt_m.group(3)) if rtt_m else 0.0,
|
||
"rtt_stddev": float(rtt_m.group(4)) if rtt_m else 0.0,
|
||
"results": results,
|
||
}
|
||
}
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – package management
|
||
# ------------------------------------------------------------------
|
||
|
||
def get_packages(self) -> list[PackageDict]:
|
||
if self._pkg_manager == "apt":
|
||
return self._get_packages_apt()
|
||
if self._pkg_manager in ("dnf", "yum"):
|
||
return self._get_packages_rpm()
|
||
if self._pkg_manager == "apk":
|
||
return self._get_packages_apk()
|
||
if self._pkg_manager == "pacman":
|
||
return self._get_packages_pacman()
|
||
raise NotImplementedError(
|
||
f"Package manager '{self._pkg_manager}' is not supported"
|
||
)
|
||
|
||
def _get_packages_apt(self) -> list[PackageDict]:
|
||
out = self._send(
|
||
"dpkg-query -W -f='${Package}\\t${Version}\\t${Installed-Size}"
|
||
"\\t${source:Package}\\t${source:Version}\\t${binary:Summary}\\n' 2>/dev/null"
|
||
)
|
||
packages: list[PackageDict] = []
|
||
for line in out.splitlines():
|
||
# Summary stays last and keeps whatever it contains: maxsplit must
|
||
# equal the number of tabs the format writes, not the field count.
|
||
parts = line.split("\t", 5)
|
||
if len(parts) < 2:
|
||
continue
|
||
name = parts[0].strip()
|
||
version = parts[1].strip()
|
||
size = int(parts[2].strip()) * 1024 if len(parts) > 2 and parts[2].strip().isdigit() else 0
|
||
# Debian source package (e.g. openssh-server → openssh) for OSV matching.
|
||
source_package = parts[3].strip() if len(parts) > 3 and parts[3].strip() else name
|
||
# And its version, which is a different number from this package's.
|
||
#
|
||
# OSV states Debian ranges in *source* versions. A source package
|
||
# that ships several binaries gives each its own upstream version:
|
||
# libldb2 is 2:2.11.0+samba4.22.11+dfsg-… while its source, samba,
|
||
# is 2:4.22.11+dfsg-…. A consumer matching on source_package and
|
||
# comparing `version` compares two unrelated numbers — dpkg reads
|
||
# ldb's 2.11.0 as older than the 2:4.17.4+dfsg-1 that fixed
|
||
# CVE-2022-44640, and a host five releases past the fix was reported
|
||
# vulnerable on four packages at once.
|
||
#
|
||
# dpkg leaves this empty when it equals `Version`; so does an older
|
||
# dpkg that does not know the field at all.
|
||
source_version = parts[4].strip() if len(parts) > 4 and parts[4].strip() else version
|
||
description = parts[5].strip() if len(parts) > 5 else ""
|
||
packages.append({
|
||
"name": name,
|
||
"version": version,
|
||
"installed": True,
|
||
"description": description,
|
||
"size": size,
|
||
"source": "apt",
|
||
"source_package": source_package,
|
||
"source_version": source_version,
|
||
})
|
||
return packages
|
||
|
||
def _get_packages_rpm(self) -> list[PackageDict]:
|
||
out = self._send(
|
||
"rpm -qa --queryformat '%{NAME}\\t%{VERSION}-%{RELEASE}\\t%{SIZE}\\t%{SUMMARY}\\n' 2>/dev/null"
|
||
)
|
||
packages: list[PackageDict] = []
|
||
for line in out.splitlines():
|
||
parts = line.split("\t", 3)
|
||
if len(parts) < 2:
|
||
continue
|
||
packages.append({
|
||
"name": parts[0].strip(),
|
||
"version": parts[1].strip(),
|
||
"installed": True,
|
||
"description": parts[3].strip() if len(parts) > 3 else "",
|
||
"size": int(parts[2].strip()) if len(parts) > 2 and parts[2].strip().isdigit() else 0,
|
||
"source": self._pkg_manager or "rpm",
|
||
})
|
||
return packages
|
||
|
||
def _get_packages_apk(self) -> list[PackageDict]:
|
||
out = self._send("apk info -v 2>/dev/null")
|
||
packages: list[PackageDict] = []
|
||
for line in out.splitlines():
|
||
# openssh-9.3_p2-r4 OpenSSH
|
||
m = re.match(r"^(\S+)-(\d[\S]*)\s*(.*)", line)
|
||
if not m:
|
||
continue
|
||
packages.append({
|
||
"name": m.group(1),
|
||
"version": m.group(2),
|
||
"installed": True,
|
||
"description": m.group(3).strip(),
|
||
"size": 0,
|
||
"source": "apk",
|
||
})
|
||
return packages
|
||
|
||
def _get_packages_pacman(self) -> list[PackageDict]:
|
||
out = self._send("pacman -Q 2>/dev/null")
|
||
packages: list[PackageDict] = []
|
||
for line in out.splitlines():
|
||
parts = line.split(None, 1)
|
||
if len(parts) < 2:
|
||
continue
|
||
packages.append({
|
||
"name": parts[0],
|
||
"version": parts[1],
|
||
"installed": True,
|
||
"description": "",
|
||
"size": 0,
|
||
"source": "pacman",
|
||
})
|
||
return packages
|
||
|
||
def search_packages(self, query: str) -> List[dict[str, Any]]:
|
||
"""Search available (installable) packages matching *query*."""
|
||
from shlex import quote as _q
|
||
safe_q = _q(query)
|
||
installed = {p["name"] for p in self.get_packages()}
|
||
packages: List[dict[str, Any]] = []
|
||
|
||
if self._pkg_manager == "apt":
|
||
out = self._send(f"apt-cache search {safe_q} 2>/dev/null")
|
||
versions: Dict[str, str] = {}
|
||
ver_out = self._send(f"apt-cache show {safe_q} 2>/dev/null | grep -E '^(Package|Version):' || true")
|
||
cur_pkg = ""
|
||
for line in ver_out.splitlines():
|
||
if line.startswith("Package:"):
|
||
cur_pkg = line.split(":", 1)[1].strip()
|
||
elif line.startswith("Version:") and cur_pkg:
|
||
versions[cur_pkg] = line.split(":", 1)[1].strip()
|
||
for line in out.splitlines():
|
||
if " - " not in line:
|
||
continue
|
||
name, _, description = line.partition(" - ")
|
||
name = name.strip()
|
||
packages.append({
|
||
"name": name,
|
||
"version": versions.get(name, ""),
|
||
"installed": name in installed,
|
||
"description": description.strip(),
|
||
"size": 0,
|
||
"source": "apt",
|
||
})
|
||
|
||
elif self._pkg_manager in ("dnf", "yum"):
|
||
cmd = "dnf" if self._pkg_manager == "dnf" else "yum"
|
||
out = self._send(f"{cmd} search {safe_q} 2>/dev/null || true")
|
||
for line in out.splitlines():
|
||
if " : " not in line:
|
||
continue
|
||
pkg_ver, _, description = line.partition(" : ")
|
||
name = pkg_ver.split(".")[0].strip()
|
||
version = ""
|
||
packages.append({
|
||
"name": name,
|
||
"version": version,
|
||
"installed": name in installed,
|
||
"description": description.strip(),
|
||
"size": 0,
|
||
"source": self._pkg_manager or "rpm",
|
||
})
|
||
|
||
elif self._pkg_manager == "apk":
|
||
out = self._send(f"apk search {safe_q} 2>/dev/null")
|
||
for line in out.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
m = re.match(r"^(.*?)-(\d\S*)(?:\s+(.*))?$", line)
|
||
if m:
|
||
name, version, description = m.group(1), m.group(2), (m.group(3) or "")
|
||
else:
|
||
name, version, description = line, "", ""
|
||
packages.append({
|
||
"name": name,
|
||
"version": version,
|
||
"installed": name in installed,
|
||
"description": description,
|
||
"size": 0,
|
||
"source": "apk",
|
||
})
|
||
|
||
elif self._pkg_manager == "pacman":
|
||
out = self._send(f"pacman -Ss {safe_q} 2>/dev/null || true")
|
||
lines = out.splitlines()
|
||
i = 0
|
||
while i < len(lines):
|
||
line = lines[i].strip()
|
||
if "/" in line and " " in line:
|
||
parts = line.split()
|
||
name_ver = parts[0].split("/")[-1] if "/" in parts[0] else parts[0]
|
||
name_parts = name_ver.rsplit(" ", 1)
|
||
name = name_parts[0]
|
||
version = parts[1] if len(parts) > 1 else ""
|
||
description = lines[i + 1].strip() if i + 1 < len(lines) else ""
|
||
packages.append({
|
||
"name": name,
|
||
"version": version,
|
||
"installed": name in installed,
|
||
"description": description,
|
||
"size": 0,
|
||
"source": "pacman",
|
||
})
|
||
i += 2
|
||
continue
|
||
i += 1
|
||
|
||
return packages
|
||
|
||
def install_package(self, name: str) -> dict[str, Any]:
|
||
"""Install a package by name. Returns ``{"success": bool, "output": str}``."""
|
||
from shlex import quote as _q
|
||
safe = _q(name)
|
||
pm = self._pkg_manager
|
||
if pm == "apt":
|
||
raw = self._sudo(f"DEBIAN_FRONTEND=noninteractive apt-get install -y {safe} 2>&1 || true")
|
||
elif pm in ("dnf", "yum"):
|
||
raw = self._sudo(f"{pm} install -y {safe} 2>&1 || true")
|
||
elif pm == "apk":
|
||
raw = self._sudo(f"apk add {safe} 2>&1 || true")
|
||
elif pm == "pacman":
|
||
raw = self._sudo(f"pacman -S --noconfirm {safe} 2>&1 || true")
|
||
else:
|
||
return {"success": False, "output": f"Unsupported package manager: {pm}"}
|
||
low = raw.lower()
|
||
success = not any(kw in low for kw in ("error:", "failed", "no packages", "not found", "unable to locate", "no match"))
|
||
return {"success": success, "output": raw.strip()}
|
||
|
||
def uninstall_package(self, name: str) -> dict[str, Any]:
|
||
"""Remove a package by name. Returns ``{"success": bool, "output": str}``."""
|
||
from shlex import quote as _q
|
||
safe = _q(name)
|
||
pm = self._pkg_manager
|
||
if pm == "apt":
|
||
raw = self._sudo(f"DEBIAN_FRONTEND=noninteractive apt-get remove -y {safe} 2>&1 || true")
|
||
elif pm in ("dnf", "yum"):
|
||
raw = self._sudo(f"{pm} remove -y {safe} 2>&1 || true")
|
||
elif pm == "apk":
|
||
raw = self._sudo(f"apk del {safe} 2>&1 || true")
|
||
elif pm == "pacman":
|
||
raw = self._sudo(f"pacman -R --noconfirm {safe} 2>&1 || true")
|
||
else:
|
||
return {"success": False, "output": f"Unsupported package manager: {pm}"}
|
||
low = raw.lower()
|
||
success = not any(kw in low for kw in ("error:", "failed", "not found", "is not installed", "no packages"))
|
||
return {"success": success, "output": raw.strip()}
|
||
|
||
def get_available_updates(self) -> list[UpdateDict]:
|
||
if self._pkg_manager == "apt":
|
||
return self._get_updates_apt()
|
||
if self._pkg_manager in ("dnf", "yum"):
|
||
return self._get_updates_rpm()
|
||
if self._pkg_manager == "apk":
|
||
return self._get_updates_apk()
|
||
if self._pkg_manager == "pacman":
|
||
return self._get_updates_pacman()
|
||
raise NotImplementedError(
|
||
f"Package manager '{self._pkg_manager}' is not supported"
|
||
)
|
||
|
||
def get_device_warnings(self) -> List[dict[str, Any]]:
|
||
"""Return warning dicts for issues detected on this device.
|
||
|
||
Currently detects:
|
||
- package updates available (uses local package cache)
|
||
- apt proxy not configured (apt systems only)
|
||
"""
|
||
warnings: List[dict[str, Any]] = []
|
||
try:
|
||
updates = self.get_available_updates()
|
||
except Exception as exc:
|
||
logger.warning("get_device_warnings: get_available_updates() failed: %s", exc)
|
||
updates = []
|
||
if updates:
|
||
warnings.append({
|
||
"code": "updates_available",
|
||
"meta": {
|
||
"count": len(updates),
|
||
"packages": [u.get("name", "") for u in updates],
|
||
},
|
||
})
|
||
if self._pkg_manager == "apt" and self._apt_proxy_url:
|
||
try:
|
||
current = self._send("cat /etc/apt/apt.conf.d/00proxy 2>/dev/null || true").strip()
|
||
if self._apt_proxy_url not in current:
|
||
warnings.append({
|
||
"code": "apt_proxy_missing",
|
||
"meta": {"expected_url": self._apt_proxy_url},
|
||
})
|
||
except Exception as exc:
|
||
logger.warning("get_device_warnings: apt proxy check failed: %s", exc)
|
||
return warnings
|
||
|
||
def _get_updates_apt(self) -> list[UpdateDict]:
|
||
# apt list --upgradable does not need root; avoid sudo so it works even
|
||
# without a configured sudo password.
|
||
out = self._send(
|
||
"LC_ALL=C apt list --upgradable 2>/dev/null | grep -v '^Listing'",
|
||
read_timeout=60,
|
||
)
|
||
# Join wrapped lines: netmiko's 80-col pseudo-TTY causes long apt lines to
|
||
# break; continuation lines start with a space.
|
||
raw_lines: List[str] = []
|
||
for line in out.splitlines():
|
||
if line.startswith(" ") and raw_lines:
|
||
raw_lines[-1] += line.strip()
|
||
else:
|
||
raw_lines.append(line)
|
||
updates: list[UpdateDict] = []
|
||
for line in raw_lines:
|
||
# openssh-server/stable 1:9.2p1-2+deb12u2 amd64 [upgradable from: 1:9.2p1-2+deb12u1]
|
||
m = re.match(
|
||
r"^(\S+)/\S+\s+(\S+)\s+\S+\s+\[upgradable from:\s+(\S+)\]", line
|
||
)
|
||
if m:
|
||
updates.append({
|
||
"name": m.group(1),
|
||
"current_version": m.group(3),
|
||
"new_version": m.group(2),
|
||
})
|
||
return updates
|
||
|
||
def _get_updates_rpm(self) -> list[UpdateDict]:
|
||
cmd = "dnf check-update --quiet 2>/dev/null" if self._pkg_manager == "dnf" else "yum check-update -q 2>/dev/null"
|
||
out = self._sudo(cmd)
|
||
updates: list[UpdateDict] = []
|
||
for line in out.splitlines():
|
||
parts = line.split()
|
||
if len(parts) >= 2 and not line.startswith(" ") and "." in parts[0]:
|
||
name_arch = parts[0]
|
||
name = name_arch.rsplit(".", 1)[0] if "." in name_arch else name_arch
|
||
updates.append({
|
||
"name": name,
|
||
"current_version": "",
|
||
"new_version": parts[1],
|
||
})
|
||
return updates
|
||
|
||
def _get_updates_apk(self) -> list[UpdateDict]:
|
||
out = self._send("apk version -l '<' 2>/dev/null")
|
||
updates: list[UpdateDict] = []
|
||
for line in out.splitlines():
|
||
# openssh-9.3_p2-r3 < 9.3_p2-r4
|
||
m = re.match(r"^(\S+)-(\S+)\s+<\s+(\S+)", line)
|
||
if m:
|
||
updates.append({
|
||
"name": m.group(1),
|
||
"current_version": m.group(2),
|
||
"new_version": m.group(3),
|
||
})
|
||
return updates
|
||
|
||
def _get_updates_pacman(self) -> list[UpdateDict]:
|
||
out = self._send("pacman -Qu 2>/dev/null")
|
||
updates: list[UpdateDict] = []
|
||
for line in out.splitlines():
|
||
# openssh 9.3p2-1 -> 9.4p1-1
|
||
m = re.match(r"^(\S+)\s+(\S+)\s+->\s+(\S+)", line)
|
||
if m:
|
||
updates.append({
|
||
"name": m.group(1),
|
||
"current_version": m.group(2),
|
||
"new_version": m.group(3),
|
||
})
|
||
return updates
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – apply updates
|
||
# ------------------------------------------------------------------
|
||
|
||
# Allowlist for package names – same pattern used by napalm-proxmox
|
||
_PKG_NAME_RE = re.compile(r'^[a-zA-Z0-9_\-\+\.]+$')
|
||
|
||
def apply_updates(self, packages: List[str]) -> ApplyUpdatesResultDict:
|
||
"""Upgrade *packages* (or all pending updates when the list is empty).
|
||
|
||
Package names are validated against ``^[a-zA-Z0-9_\\-\\+\\.]+$`` before
|
||
being passed to the package manager to prevent shell injection.
|
||
"""
|
||
for pkg in packages:
|
||
if not self._PKG_NAME_RE.match(pkg):
|
||
raise ValueError(f"Invalid package name: {pkg!r}")
|
||
|
||
if self._pkg_manager == "apt":
|
||
return self._apply_updates_apt(packages)
|
||
if self._pkg_manager in ("dnf", "yum"):
|
||
return self._apply_updates_rpm(packages)
|
||
if self._pkg_manager == "apk":
|
||
return self._apply_updates_apk(packages)
|
||
if self._pkg_manager == "pacman":
|
||
return self._apply_updates_pacman(packages)
|
||
raise NotImplementedError(
|
||
f"Package manager '{self._pkg_manager}' is not supported"
|
||
)
|
||
|
||
def _apply_updates_apt(self, packages: List[str]) -> ApplyUpdatesResultDict:
|
||
pkg_args = " ".join(packages) if packages else "--with-new-pkgs"
|
||
cmd = (
|
||
"DEBIAN_FRONTEND=noninteractive apt-get install --only-upgrade -y "
|
||
f"{pkg_args} 2>&1"
|
||
if packages else
|
||
"DEBIAN_FRONTEND=noninteractive apt-get upgrade -y 2>&1"
|
||
)
|
||
try:
|
||
output = self._sudo(cmd, read_timeout=600)
|
||
success = not re.search(r'^E:', output, re.MULTILINE)
|
||
result: ApplyUpdatesResultDict = {"success": success, "output": output}
|
||
if not success:
|
||
m = re.search(r'^E:.*', output, re.MULTILINE)
|
||
result["error"] = m.group(0) if m else "apt-get exited with errors"
|
||
return result
|
||
except Exception as exc:
|
||
return {"success": False, "output": "", "error": str(exc)}
|
||
|
||
def _apply_updates_rpm(self, packages: List[str]) -> ApplyUpdatesResultDict:
|
||
bin_ = self._pkg_manager # "dnf" or "yum"
|
||
if packages:
|
||
pkg_args = " ".join(packages)
|
||
cmd = f"{bin_} upgrade -y {pkg_args} 2>&1"
|
||
else:
|
||
cmd = f"{bin_} upgrade -y 2>&1"
|
||
try:
|
||
output = self._sudo(cmd, read_timeout=600)
|
||
# dnf/yum signal failure via "Error:" lines or non-zero exit;
|
||
# since we can't check the exit code directly, look for error markers.
|
||
success = not re.search(r'^Error:', output, re.MULTILINE | re.IGNORECASE)
|
||
result: ApplyUpdatesResultDict = {"success": success, "output": output}
|
||
if not success:
|
||
m = re.search(r'^Error:.*', output, re.MULTILINE | re.IGNORECASE)
|
||
result["error"] = m.group(0) if m else f"{bin_} exited with errors"
|
||
return result
|
||
except Exception as exc:
|
||
return {"success": False, "output": "", "error": str(exc)}
|
||
|
||
def _apply_updates_apk(self, packages: List[str]) -> ApplyUpdatesResultDict:
|
||
if packages:
|
||
pkg_args = " ".join(packages)
|
||
cmd = f"apk upgrade {pkg_args} 2>&1"
|
||
else:
|
||
cmd = "apk upgrade 2>&1"
|
||
try:
|
||
output = self._sudo(cmd, read_timeout=300)
|
||
success = "ERROR" not in output.upper().split("\n")[0] if output else True
|
||
result: ApplyUpdatesResultDict = {"success": success, "output": output}
|
||
if not success:
|
||
result["error"] = "apk upgrade reported an error"
|
||
return result
|
||
except Exception as exc:
|
||
return {"success": False, "output": "", "error": str(exc)}
|
||
|
||
def _apply_updates_pacman(self, packages: List[str]) -> ApplyUpdatesResultDict:
|
||
if packages:
|
||
pkg_args = " ".join(packages)
|
||
cmd = f"pacman --noconfirm -S {pkg_args} 2>&1"
|
||
else:
|
||
cmd = "pacman --noconfirm -Syu 2>&1"
|
||
try:
|
||
output = self._sudo(cmd, read_timeout=300)
|
||
success = "error" not in output.lower()
|
||
result: ApplyUpdatesResultDict = {"success": success, "output": output}
|
||
if not success:
|
||
result["error"] = "pacman reported an error"
|
||
return result
|
||
except Exception as exc:
|
||
return {"success": False, "output": "", "error": str(exc)}
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – services (systemd)
|
||
# ------------------------------------------------------------------
|
||
|
||
def get_services(self) -> list[ServiceDict]:
|
||
"""Return systemd service units (falls back to service --status-all on SysV)."""
|
||
out = self._send(
|
||
"systemctl list-units --type=service --all --no-legend --no-pager "
|
||
"--plain 2>/dev/null"
|
||
)
|
||
if not out:
|
||
return self._get_services_sysv()
|
||
|
||
services: list[ServiceDict] = []
|
||
for line in out.splitlines():
|
||
# ssh.service loaded active running OpenBSD Secure Shell server
|
||
parts = line.split(None, 4)
|
||
if len(parts) < 4:
|
||
continue
|
||
unit, load, active, sub = parts[0], parts[1], parts[2], parts[3]
|
||
name = unit.removesuffix(".service")
|
||
running = active == "active" and sub == "running"
|
||
enabled_out = self._send(
|
||
f"systemctl is-enabled {unit} 2>/dev/null"
|
||
)
|
||
enabled = enabled_out.strip() == "enabled"
|
||
|
||
# Retrieve main PID for running services
|
||
pid = 0
|
||
if running:
|
||
pid_out = self._send(
|
||
f"systemctl show -p MainPID --value {unit} 2>/dev/null"
|
||
)
|
||
try:
|
||
pid = int(pid_out.strip())
|
||
except ValueError:
|
||
pid = 0
|
||
|
||
services.append({
|
||
"name": name,
|
||
"running": running,
|
||
"enabled": enabled,
|
||
"pid": pid,
|
||
})
|
||
return services
|
||
|
||
def _get_services_sysv(self) -> list[ServiceDict]:
|
||
out = self._send("service --status-all 2>/dev/null")
|
||
services: list[ServiceDict] = []
|
||
for line in out.splitlines():
|
||
m = re.match(r"^\s*\[\s*([+\-?])\s*\]\s+(\S+)", line)
|
||
if not m:
|
||
continue
|
||
services.append({
|
||
"name": m.group(2),
|
||
"running": m.group(1) == "+",
|
||
"enabled": False,
|
||
"pid": 0,
|
||
})
|
||
return services
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – users
|
||
# ------------------------------------------------------------------
|
||
|
||
def get_users(self) -> list[UserDict]:
|
||
"""Return local user accounts from /etc/passwd plus supplementary groups."""
|
||
passwd_out = self._send("getent passwd 2>/dev/null || cat /etc/passwd")
|
||
groups_out = self._send("getent group 2>/dev/null || cat /etc/group")
|
||
|
||
# Build uid→[group] map from /etc/group
|
||
uid_to_groups: Dict[int, List[str]] = {}
|
||
for line in groups_out.splitlines():
|
||
parts = line.split(":")
|
||
if len(parts) < 4:
|
||
continue
|
||
gname = parts[0]
|
||
members = [m.strip() for m in parts[3].split(",") if m.strip()]
|
||
for member in members:
|
||
# We'll convert username→uid below; collect by username first
|
||
uid_to_groups.setdefault(-1, []) # placeholder
|
||
|
||
# Simpler: collect username→groups, then join with passwd
|
||
username_to_groups: Dict[str, List[str]] = {}
|
||
for line in groups_out.splitlines():
|
||
parts = line.split(":")
|
||
if len(parts) < 4:
|
||
continue
|
||
gname = parts[0]
|
||
members = [m.strip() for m in parts[3].split(",") if m.strip()]
|
||
for member in members:
|
||
username_to_groups.setdefault(member, []).append(gname)
|
||
|
||
users: list[UserDict] = []
|
||
for line in passwd_out.splitlines():
|
||
parts = line.split(":")
|
||
if len(parts) < 7:
|
||
continue
|
||
username, _, uid_s, gid_s, _, home, shell = parts[:7]
|
||
try:
|
||
uid, gid = int(uid_s), int(gid_s)
|
||
except ValueError:
|
||
continue
|
||
users.append({
|
||
"username": username,
|
||
"uid": uid,
|
||
"gid": gid,
|
||
"home": home,
|
||
"shell": shell,
|
||
"groups": username_to_groups.get(username, []),
|
||
})
|
||
return users
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – processes
|
||
# ------------------------------------------------------------------
|
||
|
||
def get_processes(self) -> list[ProcessDict]:
|
||
"""Return running processes via ``ps axo``."""
|
||
out = self._send(
|
||
"ps axo pid,ppid,user:20,pcpu,pmem,vsz,rss,tty,stat,lstart,args "
|
||
"--no-headers 2>/dev/null"
|
||
)
|
||
processes: list[ProcessDict] = []
|
||
for line in out.splitlines():
|
||
parts = line.split(None, 10)
|
||
if len(parts) < 11:
|
||
continue
|
||
try:
|
||
pid = int(parts[0])
|
||
ppid = int(parts[1])
|
||
user = parts[2]
|
||
cpu = float(parts[3])
|
||
mem = float(parts[4])
|
||
vsz = int(parts[5])
|
||
rss = int(parts[6])
|
||
tty = parts[7] if parts[7] != "?" else ""
|
||
state = parts[8][0] if parts[8] else "?"
|
||
# lstart is 5 tokens: "Mon May 27 12:34:56 2024" → parts[9..13]
|
||
# args starts at parts[14] but we merged from 10 onward
|
||
# With --no-headers and ps axo, lstart takes 5 parts
|
||
# Rebuild: parts[9] is start, args is parts[10]
|
||
started = parts[9]
|
||
command = parts[10]
|
||
except (ValueError, IndexError):
|
||
continue
|
||
processes.append({
|
||
"pid": pid,
|
||
"ppid": ppid,
|
||
"user": user,
|
||
"cpu": cpu,
|
||
"memory": mem,
|
||
"vsz": vsz,
|
||
"rss": rss,
|
||
"tty": tty,
|
||
"state": state,
|
||
"started": started,
|
||
"command": command,
|
||
})
|
||
return processes
|
||
|
||
# ------------------------------------------------------------------
|
||
# OSDriver – cron jobs
|
||
# ------------------------------------------------------------------
|
||
|
||
def get_cron_jobs(self) -> list[CronJobDict]:
|
||
"""Return cron entries from user crontabs and /etc/cron.d."""
|
||
jobs: list[CronJobDict] = []
|
||
|
||
# /etc/cron.d/* — system-wide cron fragments (include user field)
|
||
cron_d_files = self._send("ls /etc/cron.d/ 2>/dev/null").splitlines()
|
||
for fname in cron_d_files:
|
||
fname = fname.strip()
|
||
if not fname:
|
||
continue
|
||
content = self._send(f"cat /etc/cron.d/{fname} 2>/dev/null")
|
||
for line in content.splitlines():
|
||
job = self._parse_cron_line(line, source_user="root", has_user_field=True)
|
||
if job:
|
||
jobs.append(job)
|
||
|
||
# Per-user crontabs from /var/spool/cron/crontabs (Debian) or /var/spool/cron (RHEL)
|
||
for spool_dir in ("/var/spool/cron/crontabs", "/var/spool/cron"):
|
||
ls_out = self._send(f"ls {spool_dir} 2>/dev/null")
|
||
for uname in ls_out.splitlines():
|
||
uname = uname.strip()
|
||
if not uname:
|
||
continue
|
||
content = self._send(f"cat {spool_dir}/{uname} 2>/dev/null")
|
||
for line in content.splitlines():
|
||
job = self._parse_cron_line(line, source_user=uname, has_user_field=False)
|
||
if job:
|
||
jobs.append(job)
|
||
|
||
return jobs
|
||
|
||
@staticmethod
|
||
def _parse_cron_line(
|
||
line: str, source_user: str, has_user_field: bool
|
||
) -> Optional[CronJobDict]:
|
||
"""Parse a single crontab line; returns ``None`` for comments/blanks."""
|
||
stripped = line.strip()
|
||
# Remove trailing comment
|
||
comment = ""
|
||
if "#" in stripped:
|
||
idx = stripped.index("#")
|
||
comment = stripped[idx + 1:].strip()
|
||
stripped = stripped[:idx].strip()
|
||
|
||
if not stripped or stripped.startswith("@") or stripped.startswith("MAILTO"):
|
||
return None
|
||
|
||
parts = stripped.split(None, 6 if has_user_field else 5)
|
||
expected = 6 if has_user_field else 5
|
||
if len(parts) < expected:
|
||
return None
|
||
|
||
schedule = " ".join(parts[:5])
|
||
if has_user_field:
|
||
user = parts[5]
|
||
command = parts[6] if len(parts) > 6 else ""
|
||
else:
|
||
user = source_user
|
||
command = parts[5] if len(parts) > 5 else ""
|
||
|
||
job: CronJobDict = {
|
||
"user": user,
|
||
"schedule": schedule,
|
||
"command": command,
|
||
}
|
||
if comment:
|
||
job["description"] = comment
|
||
return job
|
||
|
||
# ------------------------------------------------------------------
|
||
# Docker
|
||
# ------------------------------------------------------------------
|
||
|
||
def _docker_bin(self) -> str:
|
||
"""Path to the docker binary.
|
||
|
||
A hook rather than a literal because the Docker *logic* is the same
|
||
everywhere while the *location* is not: QTS ships Container Station's
|
||
docker under /share/<pool>/.qpkg/ and never puts it on PATH. Subclasses
|
||
override this one method instead of reimplementing the surface.
|
||
"""
|
||
return "docker"
|
||
|
||
def get_docker_info(self) -> DockerInfoDict:
|
||
"""Return information about the local Docker environment.
|
||
|
||
Uses a single SSH call to collect all Docker data at once, eliminating
|
||
per-section round-trip overhead. Labels from ``docker images`` are used
|
||
directly for the OCI version field — no separate ``docker image inspect``
|
||
needed.
|
||
|
||
Returns a dict with keys:
|
||
- ``available`` (bool) — False if docker is not installed/accessible
|
||
- ``version`` (str) — Docker Engine version string
|
||
- ``containers`` (list) — list of container dicts
|
||
- ``images`` (list) — list of image dicts
|
||
- ``volumes`` (list) — list of volume dicts
|
||
- ``networks`` (list) — list of network dicts
|
||
"""
|
||
import json as _json
|
||
|
||
docker = self._docker_bin()
|
||
|
||
# Check docker binary first (docker --version doesn't need socket access)
|
||
if not self._send(f"command -v {docker} 2>/dev/null").strip():
|
||
return {"available": False}
|
||
|
||
# Verify socket access — docker ps is cheaper and fails immediately on permission errors
|
||
ps_check = self._send(f"{docker} ps 2>&1")
|
||
if "permission denied" in ps_check.lower() or "cannot connect" in ps_check.lower():
|
||
return {"available": False, "permission_denied": True}
|
||
|
||
version = self._send(f"{docker} --version 2>/dev/null").strip()
|
||
|
||
combined = self._send(
|
||
"echo '---CONTAINERS---'; "
|
||
f"{docker} ps -a --format '{{{{json .}}}}' 2>/dev/null; "
|
||
"echo '---IMAGES---'; "
|
||
f"{docker} images --format '{{{{json .}}}}' 2>/dev/null; "
|
||
"echo '---VOLUMES---'; "
|
||
f"{docker} volume ls --format '{{{{json .}}}}' 2>/dev/null; "
|
||
"echo '---NETWORKS---'; "
|
||
f"{docker} network ls --format '{{{{json .}}}}' 2>/dev/null; "
|
||
"echo '---CONFIGIMAGES---'; "
|
||
f"{docker} ps -aq 2>/dev/null | xargs -r {docker} inspect "
|
||
f"--format '{{{{.Id}}}}|{{{{.Config.Image}}}}|{{{{.Image}}}}' 2>/dev/null",
|
||
read_timeout=60,
|
||
)
|
||
|
||
if "---CONTAINERS---" not in combined:
|
||
return {"available": False}
|
||
|
||
# Split into sections
|
||
def _section(text: str, marker: str, next_marker: str) -> str:
|
||
start = text.find(marker)
|
||
if start == -1:
|
||
return ""
|
||
start += len(marker)
|
||
end = text.find(next_marker, start)
|
||
return text[start:end] if end != -1 else text[start:]
|
||
|
||
raw_containers = _section(combined, "---CONTAINERS---", "---IMAGES---")
|
||
raw_images = _section(combined, "---IMAGES---", "---VOLUMES---")
|
||
raw_volumes = _section(combined, "---VOLUMES---", "---NETWORKS---")
|
||
raw_networks = _section(combined, "---NETWORKS---", "---CONFIGIMAGES---")
|
||
raw_cfgimages = _section(combined, "---CONFIGIMAGES---", "\x00") # sentinel
|
||
|
||
def _parse_labels(raw: Any) -> Dict[str, str]:
|
||
"""Parse Docker labels — may be a dict (JSON map) or comma-sep string."""
|
||
if isinstance(raw, dict):
|
||
return {str(k): str(v) for k, v in raw.items()}
|
||
if isinstance(raw, str) and raw:
|
||
result: Dict[str, str] = {}
|
||
for part in raw.split(","):
|
||
if "=" in part:
|
||
k, _, v = part.partition("=")
|
||
result[k.strip()] = v.strip()
|
||
return result
|
||
return {}
|
||
|
||
# Containers
|
||
containers: List[dict[str, Any]] = []
|
||
for line in raw_containers.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
obj = _json.loads(line)
|
||
labels = _parse_labels(obj.get("Labels", ""))
|
||
containers.append({
|
||
"id": obj.get("ID", ""),
|
||
"name": obj.get("Names", ""),
|
||
"image": obj.get("Image", ""),
|
||
"image_version": labels.get("org.opencontainers.image.version", ""),
|
||
"command": obj.get("Command", ""),
|
||
"created": obj.get("CreatedAt", ""),
|
||
"status": obj.get("Status", ""),
|
||
"ports": obj.get("Ports", ""),
|
||
"state": obj.get("State", ""),
|
||
"compose_project": labels.get("com.docker.compose.project", ""),
|
||
"compose_service": labels.get("com.docker.compose.service", ""),
|
||
"compose_file": labels.get("com.docker.compose.project.config_files", ""),
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
# Images — OCI version comes from Labels, no separate inspect needed
|
||
images: List[dict[str, Any]] = []
|
||
for line in raw_images.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
obj = _json.loads(line)
|
||
labels = _parse_labels(obj.get("Labels", ""))
|
||
images.append({
|
||
"id": obj.get("ID", ""),
|
||
"repository": obj.get("Repository", ""),
|
||
"tag": obj.get("Tag", ""),
|
||
"size": obj.get("Size", ""),
|
||
"created": obj.get("CreatedAt", ""),
|
||
"version": labels.get("org.opencontainers.image.version", ""),
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
# Stable image reference + restart-pending detection.
|
||
#
|
||
# ``docker ps`` only reports a usable tag while that tag still resolves to
|
||
# the running image. Pull a newer image without recreating the container and
|
||
# it degrades to a bare image ID — useless as a registry reference, and the
|
||
# very state in which an update is waiting. ``.Config.Image`` is the
|
||
# reference the container was created from and never degrades.
|
||
cfg_by_cid: dict[str, tuple] = {}
|
||
for line in raw_cfgimages.splitlines():
|
||
parts = line.strip().split("|")
|
||
if len(parts) != 3 or not parts[0]:
|
||
continue
|
||
cid, cfg_ref, run_id = parts
|
||
cfg_by_cid[cid[:12]] = (cfg_ref.strip(), run_id.strip())
|
||
|
||
tag_index: dict[str, tuple] = {}
|
||
for im in images:
|
||
repo, tag = im.get("repository", ""), im.get("tag", "")
|
||
if not repo or not tag or "<none>" in (repo, tag):
|
||
continue
|
||
tag_index[f"{repo}:{tag}"] = (_short_image_id(im.get("id", "")), im.get("version", ""))
|
||
|
||
for c in containers:
|
||
cfg_ref, run_id = cfg_by_cid.get(c.get("id", "")[:12], ("", ""))
|
||
if not cfg_ref:
|
||
continue
|
||
c["image_ref"] = cfg_ref
|
||
c["running_image_id"] = _short_image_id(run_id)
|
||
c["restart_pending"] = False
|
||
c["pending_version"] = ""
|
||
# A stopped container is not "pending a restart" in any useful sense.
|
||
if c.get("state") != "running":
|
||
continue
|
||
tag_id, tag_version = tag_index.get(cfg_ref, ("", ""))
|
||
if tag_id and c["running_image_id"] and tag_id != c["running_image_id"]:
|
||
c["restart_pending"] = True
|
||
c["pending_version"] = tag_version
|
||
|
||
# Volumes
|
||
volumes: List[dict[str, Any]] = []
|
||
for line in raw_volumes.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
obj = _json.loads(line)
|
||
volumes.append({
|
||
"name": obj.get("Name", ""),
|
||
"driver": obj.get("Driver", ""),
|
||
"mountpoint": obj.get("Mountpoint", ""),
|
||
"scope": obj.get("Scope", ""),
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
# Networks
|
||
networks: List[dict[str, Any]] = []
|
||
for line in raw_networks.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
obj = _json.loads(line)
|
||
networks.append({
|
||
"id": obj.get("ID", ""),
|
||
"name": obj.get("Name", ""),
|
||
"driver": obj.get("Driver", ""),
|
||
"scope": obj.get("Scope", ""),
|
||
"ipv6": obj.get("IPv6", ""),
|
||
"internal": obj.get("Internal", ""),
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
return {
|
||
"available": True,
|
||
"version": version,
|
||
"containers": containers,
|
||
"images": images,
|
||
"volumes": volumes,
|
||
"networks": networks,
|
||
"outdated_images": [], # populated by separate check_docker_outdated task
|
||
}
|
||
|
||
def get_docker_outdated(self, containers: List[Dict]) -> List[str]:
|
||
"""Check registry for available updates for all container images.
|
||
|
||
Runs ``docker buildx imagetools inspect`` (metadata-only, no download)
|
||
for each unique image referenced by a container. Intended to be called
|
||
from a separate Celery task on a long interval (e.g. every 3 hours) so
|
||
it never blocks the main device poll.
|
||
|
||
Returns a list of image references that have a newer digest available.
|
||
"""
|
||
outdated_images: List[str] = []
|
||
# Prefer the reference the container was created from. `image` is whatever
|
||
# `docker ps` displayed, which collapses to a bare image ID once the tag has
|
||
# moved on — and an image ID is not something a registry can resolve.
|
||
candidate_images: List[str] = []
|
||
for c in containers:
|
||
ref = (c.get("image_ref") or c.get("image") or "").strip()
|
||
if not ref or "@sha256:" in ref: # skip digest-pinned
|
||
continue
|
||
if _looks_like_image_id(ref):
|
||
logger.warning(
|
||
"container %s reports image ID %r instead of a tag — cannot ask the "
|
||
"registry about it; skipping update check",
|
||
c.get("name", "?"), ref,
|
||
)
|
||
continue
|
||
if ref not in candidate_images:
|
||
candidate_images.append(ref)
|
||
for img_name in candidate_images:
|
||
try:
|
||
local_raw = self._send(
|
||
f"{self._docker_bin()} inspect {img_name!r} "
|
||
f"--format '{{{{index .RepoDigests 0}}}}' 2>/dev/null",
|
||
read_timeout=5,
|
||
).strip()
|
||
if not local_raw or "@" not in local_raw:
|
||
continue # locally built or not yet pulled
|
||
local_digest = local_raw.split("@", 1)[1]
|
||
|
||
remote_full = self._send(
|
||
f"{self._docker_bin()} buildx imagetools inspect {img_name!r} 2>&1",
|
||
read_timeout=30,
|
||
).strip()
|
||
if ("429" in remote_full
|
||
or "Too Many Requests" in remote_full
|
||
or "toomanyrequests" in remote_full):
|
||
logger.warning(
|
||
"Docker Hub rate limit hit for %s — run "
|
||
"'docker login' on the device to avoid this",
|
||
img_name,
|
||
)
|
||
continue
|
||
remote_digest = ""
|
||
for _line in remote_full.splitlines():
|
||
_ls = _line.strip()
|
||
if _ls.startswith("Digest:"):
|
||
remote_digest = _ls[7:].strip()
|
||
break
|
||
if not remote_digest or not remote_digest.startswith("sha256:"):
|
||
# Silence here is indistinguishable from "up to date" — say so.
|
||
logger.warning(
|
||
"no digest returned for %s; skipping update check. Registry said: %s",
|
||
img_name, remote_full[:200].replace("\n", " ") or "(nothing)",
|
||
)
|
||
continue
|
||
if local_digest != remote_digest:
|
||
outdated_images.append(img_name)
|
||
except Exception as exc:
|
||
logger.warning("image update check for %s: %s", img_name, exc)
|
||
return outdated_images
|
||
|
||
def reconstruct_docker_run(self, container_id: str) -> dict | None:
|
||
"""Return the information needed to recreate a standalone container.
|
||
|
||
Parses ``docker inspect`` JSON and returns a dict with:
|
||
- ``name`` — container name (without leading slash)
|
||
- ``image`` — current image reference
|
||
- ``run_args`` — list of CLI args for ``docker run`` (without image/cmd)
|
||
- ``cmd`` — command override (may be empty list)
|
||
- ``entrypoint`` — entrypoint override (may be empty list)
|
||
|
||
Returns None if the container does not exist or inspect fails.
|
||
"""
|
||
import json as _json
|
||
import shlex as _shlex
|
||
|
||
raw = self._send(
|
||
f"{self._docker_bin()} inspect {_shlex.quote(container_id)} 2>/dev/null",
|
||
read_timeout=10,
|
||
).strip()
|
||
if not raw:
|
||
return None
|
||
try:
|
||
data = _json.loads(raw)
|
||
except Exception:
|
||
return None
|
||
if not data:
|
||
return None
|
||
c = data[0]
|
||
|
||
name = c.get("Name", "").lstrip("/")
|
||
cfg = c.get("Config", {})
|
||
hcfg = c.get("HostConfig", {})
|
||
net_settings = c.get("NetworkSettings", {})
|
||
|
||
args: List[str] = ["--name", name]
|
||
|
||
# Restart policy
|
||
rp = hcfg.get("RestartPolicy", {})
|
||
rp_name = rp.get("Name", "no")
|
||
if rp_name and rp_name != "no":
|
||
max_retry = rp.get("MaximumRetryCount", 0)
|
||
if rp_name == "on-failure" and max_retry:
|
||
args += ["--restart", f"on-failure:{max_retry}"]
|
||
else:
|
||
args += ["--restart", rp_name]
|
||
|
||
# Hostname
|
||
hostname = cfg.get("Hostname", "")
|
||
if hostname and hostname != name[:12]:
|
||
args += ["--hostname", hostname]
|
||
|
||
# Environment (skip vars that look like Docker-injected metadata)
|
||
_skip_prefixes = ("PATH=", "HOME=", "TERM=", "HOSTNAME=")
|
||
for env in cfg.get("Env") or []:
|
||
if not any(env.startswith(p) for p in _skip_prefixes):
|
||
args += ["-e", env]
|
||
|
||
# Volume binds
|
||
for bind in hcfg.get("Binds") or []:
|
||
args += ["-v", bind]
|
||
|
||
# Port bindings
|
||
for container_port, host_bindings in (hcfg.get("PortBindings") or {}).items():
|
||
for hb in (host_bindings or []):
|
||
host_ip = hb.get("HostIp", "")
|
||
host_port = hb.get("HostPort", "")
|
||
if host_ip:
|
||
args += ["-p", f"{host_ip}:{host_port}:{container_port}"]
|
||
else:
|
||
args += ["-p", f"{host_port}:{container_port}"]
|
||
|
||
# Network mode
|
||
net_mode = hcfg.get("NetworkMode", "default")
|
||
if net_mode not in ("default", "bridge"):
|
||
args += ["--network", net_mode]
|
||
else:
|
||
# Check for custom networks from NetworkSettings
|
||
for net_name in (net_settings.get("Networks") or {}):
|
||
if net_name not in ("bridge", "host", "none"):
|
||
args += ["--network", net_name]
|
||
break
|
||
|
||
# Privileged
|
||
if hcfg.get("Privileged"):
|
||
args.append("--privileged")
|
||
|
||
# Cap-add
|
||
for cap in hcfg.get("CapAdd") or []:
|
||
args += ["--cap-add", cap]
|
||
|
||
# Devices
|
||
for dev in hcfg.get("Devices") or []:
|
||
host_p = dev.get("PathOnHost", "")
|
||
ctr_p = dev.get("PathInContainer", "")
|
||
perms = dev.get("CgroupPermissions", "rwm")
|
||
if host_p:
|
||
args += ["--device", f"{host_p}:{ctr_p}:{perms}"]
|
||
|
||
# Extra hosts
|
||
for eh in hcfg.get("ExtraHosts") or []:
|
||
args += ["--add-host", eh]
|
||
|
||
# DNS
|
||
for dns in hcfg.get("Dns") or []:
|
||
args += ["--dns", dns]
|
||
|
||
# Labels (skip Docker-internal labels)
|
||
_skip_label_prefixes = ("com.docker.compose.", "org.opencontainers.")
|
||
for k, v in (cfg.get("Labels") or {}).items():
|
||
if not any(k.startswith(p) for p in _skip_label_prefixes):
|
||
args += ["--label", f"{k}={v}"]
|
||
|
||
# Detach always
|
||
args.append("-d")
|
||
|
||
return {
|
||
"name": name,
|
||
"image": cfg.get("Image", ""),
|
||
"run_args": args,
|
||
"cmd": cfg.get("Cmd") or [],
|
||
"entrypoint": cfg.get("Entrypoint") or [],
|
||
}
|
||
|
||
# ── Device actions ────────────────────────────────────────────────────────
|
||
|
||
def get_snmp_config(self) -> Optional[SNMPConfigDict]:
|
||
"""Return SNMP agent config if snmpd is installed and running."""
|
||
try:
|
||
running = (
|
||
self._send("systemctl is-active snmpd 2>/dev/null || true").strip()
|
||
== "active"
|
||
)
|
||
if not running:
|
||
return None
|
||
|
||
# Parse community string from snmpd.conf
|
||
community = "public"
|
||
port = 161
|
||
try:
|
||
conf = self._send(
|
||
"grep -E '^[[:space:]]*(ro|rw)?community' /etc/snmp/snmpd.conf 2>/dev/null"
|
||
" | head -5"
|
||
)
|
||
for line in conf.splitlines():
|
||
parts = line.split()
|
||
if not parts:
|
||
continue
|
||
kw = parts[0].lower()
|
||
if kw in ("rocommunity", "rwcommunity", "rocommunity6", "rwcommunity6"):
|
||
if len(parts) >= 2:
|
||
community = parts[1]
|
||
break
|
||
elif kw == "com2sec" and len(parts) >= 4:
|
||
# com2sec notConfigUser default <community>
|
||
community = parts[3]
|
||
break
|
||
except Exception:
|
||
pass
|
||
|
||
# Detect port override
|
||
try:
|
||
port_line = self._send(
|
||
"grep -E '^agentAddress' /etc/snmp/snmpd.conf 2>/dev/null | head -1"
|
||
).strip()
|
||
if port_line:
|
||
m = re.search(r':(\d+)', port_line)
|
||
if m:
|
||
port = int(m.group(1))
|
||
except Exception:
|
||
pass
|
||
|
||
return SNMPConfigDict(running=True, community=community, port=port, version="2c")
|
||
except Exception as exc:
|
||
logger.debug("get_snmp_config() failed: %s", exc)
|
||
return None
|
||
|
||
def run_device_action(self, action: str) -> DeviceActionResultDict:
|
||
"""Execute a named action on the device."""
|
||
if action == "fix_docker_permissions":
|
||
return self._action_fix_docker_permissions()
|
||
if action == "fix_snmp":
|
||
return self._action_fix_snmp()
|
||
if action == "fix_apt_proxy":
|
||
return self._action_fix_apt_proxy()
|
||
if action == "apt_update_upgrade":
|
||
return self._action_apt_update_upgrade()
|
||
raise NotImplementedError(f"Unknown action: {action!r}")
|
||
|
||
def _action_apt_update_upgrade(self) -> DeviceActionResultDict:
|
||
"""Refresh the apt cache and fully upgrade all packages (apt-based systems only).
|
||
|
||
Uses full-upgrade (not plain upgrade) — plain "apt-get upgrade" refuses
|
||
to install/remove packages even when required to satisfy a newer
|
||
version's dependencies, silently leaving those updates pending.
|
||
"""
|
||
if self._pkg_manager != "apt":
|
||
return {
|
||
"success": True,
|
||
"output": f"Skipped — package manager is {self._pkg_manager!r}, not apt.",
|
||
}
|
||
|
||
sudo_check = self._send("sudo -n true 2>&1 || echo __SUDO_NEEDS_PW__")
|
||
if "__SUDO_NEEDS_PW__" in sudo_check or "password is required" in sudo_check.lower():
|
||
if not self._sudo_password:
|
||
return {
|
||
"success": False,
|
||
"output": (
|
||
"sudo requires a password on this device but none is configured in "
|
||
"netOrk. Please add the sudo password to a Credential Profile assigned "
|
||
"to this device, or configure passwordless sudo (NOPASSWD) for this user."
|
||
),
|
||
}
|
||
|
||
lines: list[str] = []
|
||
try:
|
||
out = self._sudo("apt-get update -y 2>&1", read_timeout=90)
|
||
lines.append(f"[update] {out.strip()[-300:]}")
|
||
out = self._sudo(
|
||
"DEBIAN_FRONTEND=noninteractive apt-get full-upgrade -y 2>&1", read_timeout=240
|
||
)
|
||
lines.append(f"[upgrade] {out.strip()[-300:]}")
|
||
return {"success": True, "output": "\n".join(lines)}
|
||
except Exception as exc:
|
||
lines.append(f"[error] {exc}")
|
||
return {"success": False, "output": "\n".join(lines)}
|
||
|
||
def _action_fix_snmp(self) -> DeviceActionResultDict:
|
||
"""Install, configure and start snmpd with community 'public'."""
|
||
lines: list[str] = []
|
||
|
||
# 0. Verify sudo access before attempting anything
|
||
sudo_check = self._send("sudo -n true 2>&1 || echo __SUDO_NEEDS_PW__")
|
||
if "__SUDO_NEEDS_PW__" in sudo_check or "password is required" in sudo_check.lower():
|
||
if not self._sudo_password:
|
||
return {
|
||
"success": False,
|
||
"output": (
|
||
"sudo requires a password on this device but none is configured in netOrk. "
|
||
"Please add the sudo password to a Credential Profile assigned to this device, "
|
||
"or configure passwordless sudo (NOPASSWD) for this user."
|
||
),
|
||
}
|
||
|
||
# 1. Install snmpd if missing
|
||
pkg_mgr = self._detect_pkg_manager()
|
||
if not pkg_mgr:
|
||
return {"success": False, "output": "Package manager not detected — cannot install snmpd."}
|
||
|
||
# Refresh the package index first — a freshly provisioned (or simply
|
||
# long-untouched) system's cache can be stale/empty, which makes the
|
||
# install below fail outright rather than just being slow.
|
||
if pkg_mgr == "apt":
|
||
try:
|
||
update_out = self._sudo("apt-get update -y 2>&1", read_timeout=90)
|
||
lines.append(f"[update] {update_out.strip()[-200:]}")
|
||
except Exception as exc:
|
||
lines.append(f"[warn] apt-get update failed: {exc}")
|
||
|
||
# Install both snmpd (daemon) and snmp (client tools incl. snmpget for probing)
|
||
install_cmd: dict[str, str] = {
|
||
"apt": "DEBIAN_FRONTEND=noninteractive apt-get install -y snmpd snmp 2>&1",
|
||
"dnf": "dnf install -y net-snmp net-snmp-utils 2>&1",
|
||
"yum": "yum install -y net-snmp net-snmp-utils 2>&1",
|
||
"apk": "apk add --no-cache net-snmp net-snmp-tools 2>&1",
|
||
"pacman": "pacman -Sy --noconfirm net-snmp 2>&1",
|
||
}
|
||
cmd = install_cmd.get(pkg_mgr)
|
||
if cmd:
|
||
try:
|
||
out = self._sudo(cmd, read_timeout=120)
|
||
lines.append(f"[install] {out.strip()[-200:]}")
|
||
except Exception as exc:
|
||
lines.append(f"[error] install failed: {exc}")
|
||
return {"success": False, "output": "\n".join(lines)}
|
||
|
||
# 2. Determine the IP netOrk is connecting from by checking the established SSH connection
|
||
netork_ip = ""
|
||
try:
|
||
# ss shows the remote peer of the current SSH connection
|
||
raw = self._send(
|
||
"ss -tnp 2>/dev/null | awk '/sshd/{print $5}' | head -1 | cut -d: -f1"
|
||
).strip()
|
||
if raw and raw not in ("", "0.0.0.0", "::", "127.0.0.1"):
|
||
netork_ip = raw
|
||
except Exception:
|
||
pass
|
||
|
||
# Write snmpd.conf:
|
||
# 1. Write to /tmp (no sudo needed, avoids stdin conflict with sudo -S)
|
||
# 2. sudo mv to /etc/snmp/snmpd.conf
|
||
# agentAddress udp:161 overrides Debian's localhost-only default.
|
||
import base64 as _b64
|
||
conf_str = (
|
||
"agentAddress udp:161\n"
|
||
"rocommunity public\n"
|
||
"sysLocation Managed by netOrk\n"
|
||
"sysContact netork@localhost\n"
|
||
)
|
||
conf_b64 = _b64.b64encode(conf_str.encode()).decode()
|
||
self._send(f"echo {conf_b64} | base64 -d > /tmp/netork_snmpd.conf")
|
||
self._sudo("mv /tmp/netork_snmpd.conf /etc/snmp/snmpd.conf && chown root:root /etc/snmp/snmpd.conf && chmod 644 /etc/snmp/snmpd.conf")
|
||
verify = self._send("cat /etc/snmp/snmpd.conf 2>/dev/null").strip()
|
||
if "agentAddress" in verify and "rocommunity" in verify:
|
||
lines.append("[config] Wrote /etc/snmp/snmpd.conf — agentAddress udp:161, rocommunity public.")
|
||
else:
|
||
lines.append(f"[warn] snmpd.conf write may have failed: {verify[:100]}")
|
||
|
||
# 3. Open firewall for SNMP (UDP 161) — restrict to netOrk's source IP
|
||
if netork_ip:
|
||
try:
|
||
ufw = self._send("command -v ufw 2>/dev/null").strip()
|
||
ipt = self._send("command -v iptables 2>/dev/null").strip()
|
||
if ufw:
|
||
# Expand to /24 so all containers in the same Docker network can probe
|
||
parts = netork_ip.rsplit(".", 1)
|
||
subnet = f"{parts[0]}.0/24" if len(parts) == 2 else netork_ip
|
||
fw_out = self._sudo(
|
||
f"ufw allow from {subnet} to any port 161 proto udp 2>&1", read_timeout=10
|
||
)
|
||
lines.append(f"[firewall/ufw] {fw_out.strip()[:200]}")
|
||
elif ipt:
|
||
fw_out = self._sudo(
|
||
f"iptables -C INPUT -s {netork_ip} -p udp --dport 161 -j ACCEPT 2>/dev/null"
|
||
f" || iptables -I INPUT -s {netork_ip} -p udp --dport 161 -j ACCEPT",
|
||
read_timeout=10,
|
||
)
|
||
lines.append(f"[firewall/iptables] rule added for {netork_ip}:161/udp")
|
||
except Exception as exc:
|
||
lines.append(f"[firewall] skipped — {exc}")
|
||
|
||
# 4. Restart snmpd.
|
||
# - Redirect all output to /dev/null so netmiko's prompt detection is
|
||
# never confused by service status messages.
|
||
# - Append "; echo __OK__" so there is always a known token to wait for.
|
||
import time as _time
|
||
# Stop any running snmpd (systemctl-managed or apt-started orphan)
|
||
self._sudo("systemctl stop snmpd >/dev/null 2>&1; echo s1", read_timeout=15)
|
||
self._sudo("pkill -9 snmpd >/dev/null 2>&1; echo s2", read_timeout=10)
|
||
_time.sleep(2)
|
||
# Enable and start fresh
|
||
self._sudo("systemctl enable snmpd >/dev/null 2>&1; echo s3", read_timeout=15)
|
||
self._sudo("systemctl start snmpd >/dev/null 2>&1; echo s4", read_timeout=20)
|
||
_time.sleep(2)
|
||
lines.append("[service] snmpd restarted.")
|
||
|
||
# 5. Verify snmpd responds via local SNMP probe (sysDescr.0).
|
||
# Success requires actual SNMP data types in the output, not just
|
||
# the absence of error keywords.
|
||
_time.sleep(2)
|
||
probe_out = self._send(
|
||
"snmpget -v2c -cpublic -t2 -r0 -Ov 127.0.0.1 1.3.6.1.2.1.1.1.0 2>&1 || true"
|
||
).strip()
|
||
# snmpget returns lines like "STRING: Linux ..." or "Timeticks: (n) ..."
|
||
_snmp_types = ("STRING:", "INTEGER:", "OID:", "Timeticks:", "Hex-STRING:", "IpAddress:")
|
||
success = any(t in probe_out for t in _snmp_types)
|
||
if success:
|
||
lines.append(f"[ok] SNMP probe successful — community 'public' is working.")
|
||
else:
|
||
lines.append(f"[warn] SNMP probe failed — output: {probe_out[:200]}")
|
||
|
||
return {"success": success, "output": "\n".join(lines)}
|
||
|
||
def _action_fix_docker_permissions(self) -> dict[str, Any]:
|
||
"""Add the SSH user to the 'docker' group via sudo usermod."""
|
||
user = self._send("whoami 2>/dev/null || id -un").strip().splitlines()[-1].strip()
|
||
out = self._sudo(f"usermod -aG docker {user}")
|
||
low = out.lower()
|
||
success = not any(kw in low for kw in ("error", "invalid", "no such", "command not found"))
|
||
if not out.strip():
|
||
out = f"Added {user!r} to the docker group. Reconnect or run a new poll to verify."
|
||
return {"success": success, "output": out}
|
||
|
||
def _action_fix_apt_proxy(self) -> DeviceActionResultDict:
|
||
"""Write /etc/apt/apt.conf.d/00proxy with the configured proxy URL."""
|
||
import base64 as _b64
|
||
proxy_url = self._apt_proxy_url
|
||
if not proxy_url:
|
||
return {"success": False, "output": "No apt_proxy_url configured."}
|
||
if self._pkg_manager != "apt":
|
||
return {"success": False, "output": f"Package manager is {self._pkg_manager!r}, not apt — skipping."}
|
||
|
||
content = f'Acquire::http::Proxy "{proxy_url}";\n'
|
||
content_b64 = _b64.b64encode(content.encode()).decode()
|
||
self._send(f"echo {content_b64} | base64 -d > /tmp/netork_00proxy")
|
||
self._sudo(
|
||
"mv /tmp/netork_00proxy /etc/apt/apt.conf.d/00proxy && "
|
||
"chown root:root /etc/apt/apt.conf.d/00proxy && "
|
||
"chmod 644 /etc/apt/apt.conf.d/00proxy"
|
||
)
|
||
verify = self._send("cat /etc/apt/apt.conf.d/00proxy 2>/dev/null").strip()
|
||
success = proxy_url in verify
|
||
if success:
|
||
return {"success": True, "output": f"Wrote /etc/apt/apt.conf.d/00proxy — proxy: {proxy_url}"}
|
||
return {"success": False, "output": f"Write may have failed. File content: {verify[:200]}"}
|
||
|
||
def get_vpn_tunnels(self) -> dict[str, Any]:
|
||
"""Return WireGuard status via ``wg show all dump`` (requires root/sudo).
|
||
|
||
Falls back to interface-level data from ``ip link`` + ``/proc/net/dev``
|
||
when root access is unavailable.
|
||
|
||
Full data keyed by ``wireguard-<iface>-<pubkey[:8]>`` (one entry per peer).
|
||
Fallback keyed by ``wireguard-<iface>`` (one entry per WireGuard interface).
|
||
"""
|
||
import time as _time
|
||
|
||
tunnels: dict[str, Any] = {}
|
||
|
||
# ── Attempt 1: wg show all dump via sudo ─────────────────────────────
|
||
# Write output to a fixed temp file to preserve literal tab characters.
|
||
# PTY output processing expands tabs to spaces, breaking split("\t").
|
||
import hashlib as _hashlib
|
||
_tmp = "/tmp/.netork_wg_" + _hashlib.md5(self.hostname.encode()).hexdigest()[:8]
|
||
if self._sudo_password:
|
||
self._sudo(f"wg show all dump > {_tmp} 2>&1 || true")
|
||
else:
|
||
self._send(f"sudo -n wg show all dump > {_tmp} 2>&1 || true")
|
||
raw = self._send(f"cat {_tmp} 2>/dev/null || true; rm -f {_tmp}").strip()
|
||
|
||
_perm_errors = (
|
||
"operation not permitted", "permission denied",
|
||
"a password is required", "a terminal is required",
|
||
"command not found", "not found", "no such file",
|
||
)
|
||
has_peer_data = raw and not any(e in raw.lower() for e in _perm_errors)
|
||
|
||
if has_peer_data:
|
||
# Build listen-port map from interface lines (5 tab-separated fields)
|
||
listen_ports: Dict[str, str] = {}
|
||
for line in raw.splitlines():
|
||
parts = line.split("\t")
|
||
if len(parts) == 5:
|
||
iface, _priv, _pub, port, _fwmark = parts
|
||
listen_ports[iface.strip()] = port.strip()
|
||
|
||
now = int(_time.time())
|
||
for line in raw.splitlines():
|
||
parts = line.split("\t")
|
||
if len(parts) != 9:
|
||
continue
|
||
(iface, pubkey, _psk, endpoint, allowed_ips,
|
||
latest_hs, rx_bytes, tx_bytes, _keepalive) = parts
|
||
|
||
iface = iface.strip()
|
||
pubkey = pubkey.strip()
|
||
endpoint = endpoint.strip()
|
||
|
||
remote_ip = ""
|
||
if endpoint and endpoint != "(none)":
|
||
remote_ip = endpoint.rsplit(":", 1)[0].strip("[]")
|
||
|
||
try:
|
||
hs_ts = int(latest_hs)
|
||
except (ValueError, TypeError):
|
||
hs_ts = 0
|
||
|
||
# "up" if last handshake within 3 min (WireGuard re-handshake every 2 min)
|
||
is_up = hs_ts > 0 and (now - hs_ts) < 180
|
||
|
||
try:
|
||
bytes_in = int(rx_bytes)
|
||
except (ValueError, TypeError):
|
||
bytes_in = 0
|
||
try:
|
||
bytes_out = int(tx_bytes)
|
||
except (ValueError, TypeError):
|
||
bytes_out = 0
|
||
|
||
local_port = listen_ports.get(iface, "")
|
||
local_ep = f":{local_port}" if local_port and local_port != "0" else ""
|
||
|
||
key = f"wireguard-{iface}-{pubkey[:8]}"
|
||
tunnels[key] = {
|
||
"type": "WireGuard",
|
||
"local_endpoint": local_ep,
|
||
"remote_endpoint": remote_ip,
|
||
"is_up": is_up,
|
||
"uptime": (now - hs_ts) if hs_ts > 0 else 0,
|
||
"bytes_in": bytes_in,
|
||
"bytes_out": bytes_out,
|
||
"description": f"{iface} — peer {pubkey[:16]}…",
|
||
"public_key": pubkey,
|
||
"allowed_ips": allowed_ips.strip(),
|
||
"interface": iface,
|
||
}
|
||
return tunnels
|
||
|
||
# ── Fallback: interface-level data without root ──────────────────────
|
||
# ip -j link show type wireguard → list of WireGuard interface objects
|
||
ip_raw = self._send("ip -j link show type wireguard 2>/dev/null || true").strip()
|
||
if not ip_raw or ip_raw.startswith("[") is False:
|
||
# Try stripping shell noise before the JSON
|
||
start = ip_raw.find("[")
|
||
ip_raw = ip_raw[start:] if start != -1 else ""
|
||
|
||
if not ip_raw:
|
||
return tunnels
|
||
|
||
try:
|
||
import json as _json
|
||
iface_list = _json.loads(ip_raw)
|
||
except Exception:
|
||
return tunnels
|
||
|
||
# /proc/net/dev for total RX/TX bytes per interface
|
||
proc_dev = self._send("cat /proc/net/dev 2>/dev/null || true")
|
||
proc_bytes: Dict[str, tuple] = {}
|
||
for line in proc_dev.splitlines()[2:]:
|
||
line = line.strip()
|
||
if ":" not in line:
|
||
continue
|
||
iface_name, rest = line.split(":", 1)
|
||
fields = rest.split()
|
||
try:
|
||
proc_bytes[iface_name.strip()] = (int(fields[0]), int(fields[8]))
|
||
except (IndexError, ValueError):
|
||
pass
|
||
|
||
for iface_obj in iface_list:
|
||
iface = iface_obj.get("ifname", "")
|
||
if not iface:
|
||
continue
|
||
flags = iface_obj.get("flags", [])
|
||
is_up = "UP" in flags and "LOWER_UP" in flags
|
||
rx, tx = proc_bytes.get(iface, (0, 0))
|
||
tunnels[f"wireguard-{iface}"] = {
|
||
"type": "WireGuard",
|
||
"local_endpoint": "",
|
||
"remote_endpoint": "",
|
||
"is_up": is_up,
|
||
"uptime": 0,
|
||
"bytes_in": rx,
|
||
"bytes_out": tx,
|
||
"description": f"{iface} (peer data requires root/sudo)",
|
||
"interface": iface,
|
||
}
|
||
|
||
return tunnels
|