uninstall_package judged success by searching apt/dnf/apk/pacman output
for failure words. That is guesswork in both directions: apt's commonest
failure ("E: Sub-process /usr/bin/dpkg returned an error code (1)") read
as success until the previous change, and a prerm that prints "Failed to
stop ..." while the removal completes still reads as failure. The exit
status is the answer the package manager actually gives, but every
command went through `_sudo(... || true)`, which throws it away.
Add `_sudo_status()`, which runs the command via `_sudo` followed by
`; echo __NETORK_RC=$?` and returns `(output, exit_status)` with the
marker stripped. The `|| true` of other `_sudo` callers is untouched:
they still want output rather than a status. The marker is matched only
on a line of its own with digits, so an echoed command line (literal
`$?`) is never mistaken for it. If the marker never arrives the status
is None -- unknown, not success.
uninstall_package and its dpkg fallback now use it, and
`_uninstall_failed(output, rc)` lets rc decide whenever it is known,
falling back to the keyword check only when it is not.
Behaviour change worth knowing: removing a package that is not installed
exits 0 on apt (and dnf), so it now reports success where the keyword
"is not installed" used to report failure. The package is absent
afterwards, which is what the caller asked for, and netOrk dropping it
from the installed record is then correct.
Refs christianmanivong/netork#267
2437 lines
100 KiB
Python
2437 lines
100 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"]
|
||
|
||
#: Printed after a command by ``_sudo_status`` so its exit status survives the
|
||
#: trip through an interactive shell. Matched only on a line of its own with a
|
||
#: number after it — an echoed command line carries the literal ``$?`` instead.
|
||
_RC_MARKER = "__NETORK_RC="
|
||
_RC_MARKER_RE = re.compile(rf"^{_RC_MARKER}(\d+)\s*$", re.MULTILINE)
|
||
|
||
# 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 _sudo_status(self, command: str, read_timeout: float = 100) -> tuple[str, int | None]:
|
||
"""Run *command* via sudo and return ``(output, exit_status)``.
|
||
|
||
``_sudo`` callers append ``|| true`` so a failing command yields output
|
||
instead of an error, which throws the exit status away. This variant
|
||
echoes ``$?`` straight after the sudo pipeline instead — sudo passes
|
||
the command's status through, and a failed password is non-zero too.
|
||
|
||
The status is ``None`` when the marker never arrived (output cut short),
|
||
so a caller can tell "unknown" from "succeeded".
|
||
"""
|
||
raw = self._sudo(f"{command}; echo {_RC_MARKER}$?", read_timeout=read_timeout)
|
||
matches = list(_RC_MARKER_RE.finditer(raw))
|
||
if not matches:
|
||
return raw, None
|
||
last = matches[-1]
|
||
output = (raw[: last.start()] + raw[last.end():]).strip()
|
||
return output, int(last.group(1))
|
||
|
||
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()}
|
||
|
||
#: Words in a package manager's output that mean it did not do the job.
|
||
#: Only consulted when the exit status is unknown; see ``_uninstall_failed``.
|
||
_UNINSTALL_FAILED = ("error:", "failed", "not found", "is not installed", "no packages")
|
||
|
||
def uninstall_package(self, name: str, purge: bool = False) -> dict[str, Any]:
|
||
"""Remove a package by name. Returns ``{"success": bool, "output": str}``.
|
||
|
||
``purge`` also removes the package's configuration where the package
|
||
manager distinguishes the two. Off by default: configuration somebody
|
||
may want back is not this function's to delete unless it was asked for.
|
||
|
||
It matters for more than tidiness. A package's apt source survives a
|
||
plain ``remove``, so the repository keeps being fetched on every
|
||
``apt-get update`` long after the package itself is gone — which is what
|
||
the Wazuh agent left behind on thirteen hosts.
|
||
|
||
**The dpkg fallback.** A package whose ``postinst`` failed sits at
|
||
``install ok unpacked``, and apt cannot remove it: it configures a
|
||
package before removing it, and configuring is precisely what is broken.
|
||
Seven of those thirteen hosts were in that state after an upgrade whose
|
||
postinst could not reach a manager that had been decommissioned, and on
|
||
one of them only ``dpkg --purge --force-all`` got it out.
|
||
|
||
So the fallback runs **only after apt has failed**, never as a routine
|
||
second step: forcing dpkg past its own consistency checks is a bigger
|
||
hammer than apt, and a caller who reaches for it every time will
|
||
eventually break something apt would have refused to.
|
||
"""
|
||
from shlex import quote as _q
|
||
safe = _q(name)
|
||
pm = self._pkg_manager
|
||
if pm == "apt":
|
||
action = "purge" if purge else "remove"
|
||
cmd = f"DEBIAN_FRONTEND=noninteractive apt-get {action} -y {safe} 2>&1"
|
||
elif pm in ("dnf", "yum"):
|
||
cmd = f"{pm} remove -y {safe} 2>&1"
|
||
elif pm == "apk":
|
||
# apk and pacman have no separate purge; asking for one is not an
|
||
# error, it simply has nothing extra to do.
|
||
cmd = f"apk del {safe} 2>&1"
|
||
elif pm == "pacman":
|
||
cmd = f"pacman -R --noconfirm {safe} 2>&1"
|
||
else:
|
||
return {"success": False, "output": f"Unsupported package manager: {pm}"}
|
||
|
||
raw, rc = self._sudo_status(cmd)
|
||
failed = self._uninstall_failed(raw, rc)
|
||
|
||
if failed and pm == "apt":
|
||
forced, forced_rc = self._sudo_status(f"dpkg --purge --force-all {safe} 2>&1")
|
||
raw = f"{raw.strip()}\n--- dpkg --purge --force-all ---\n{forced.strip()}"
|
||
failed = self._uninstall_failed(forced, forced_rc)
|
||
|
||
return {"success": not failed, "output": raw.strip()}
|
||
|
||
def _uninstall_failed(self, output: str, rc: int | None = None) -> bool:
|
||
"""Whether the package manager did not do the job.
|
||
|
||
The exit status decides whenever there is one (netork#267): it is the
|
||
answer the package manager actually gives, where the output is prose
|
||
that every tool phrases differently. A prerm printing "Failed to stop
|
||
…" while the removal completes is a success; a non-zero exit with
|
||
nothing alarming in the output is not.
|
||
|
||
Only when the status is unknown (``rc is None``) is the output read,
|
||
as the best answer left. apt prefixes its own errors with ``E: `` at
|
||
the start of a line, and the commonest of them — ``E: Sub-process
|
||
/usr/bin/dpkg returned an error code (1)`` — contains neither "error:"
|
||
nor "failed". The keyword list alone therefore read a failed removal
|
||
as a success, which is the worst direction for this particular answer
|
||
to be wrong in.
|
||
|
||
Matched at line start rather than anywhere: "note: " ends in "e: ".
|
||
"""
|
||
if rc is not None:
|
||
return rc != 0
|
||
low = output.lower()
|
||
if any(line.lstrip().startswith("e: ") for line in low.splitlines()):
|
||
return True
|
||
return any(kw in low for kw in self._UNINSTALL_FAILED)
|
||
|
||
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
|