Mirrors the robust fallback pattern from napalm-opnsense so the driver works regardless of which key name the caller passes in optional_args. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2380 lines
92 KiB
Python
2380 lines
92 KiB
Python
"""NAPALM driver for Proxmox VE.
|
||
|
||
Supports:
|
||
- Classic Linux networking (/etc/network/interfaces via Proxmox API)
|
||
- Software-Defined Networking (SDN): zones, VNets, subnets
|
||
- Open vSwitch (OVS) bridges, bonds, and internal ports
|
||
|
||
Connection is made via the Proxmox REST API (``proxmoxer`` library).
|
||
The driver targets the *node* level: each Proxmox node is treated as a
|
||
network device. Cluster-wide SDN information is also exposed where the
|
||
NAPALM API allows it.
|
||
|
||
Optional args
|
||
-------------
|
||
verify_ssl : bool
|
||
Verify TLS certificates (default: True).
|
||
port : int
|
||
Proxmox API port (default: 8006).
|
||
node : str
|
||
Override the target node name (default: auto-detected from hostname).
|
||
realm : str
|
||
PAM realm (default: ``pam``).
|
||
token_name : str
|
||
API token name (e.g. ``napalm@pam!mytoken``).
|
||
token_value : str
|
||
API token secret. When both token_name and token_value are provided,
|
||
token-based auth is used instead of password auth.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import re
|
||
import socket
|
||
from typing import Any
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
from napalm_device_types import HypervisorDriver
|
||
from napalm.base.exceptions import (
|
||
ConnectionException,
|
||
SessionLockedException,
|
||
)
|
||
from napalm.base.helpers import mac as napalm_mac
|
||
from napalm.base.netmiko_helpers import netmiko_args
|
||
import napalm.base.constants as C
|
||
|
||
try:
|
||
from proxmoxer import ProxmoxAPI
|
||
from proxmoxer.core import ResourceException
|
||
except ImportError as exc: # pragma: no cover
|
||
raise ImportError(
|
||
"proxmoxer is required: pip install proxmoxer"
|
||
) from exc
|
||
|
||
from napalm_proxmox import utils
|
||
|
||
# --------------------------------------------------------------------------- #
|
||
# Type aliases
|
||
# --------------------------------------------------------------------------- #
|
||
_JsonDict = dict[str, Any]
|
||
|
||
# --------------------------------------------------------------------------- #
|
||
# Driver
|
||
# --------------------------------------------------------------------------- #
|
||
|
||
|
||
class ProxmoxDriver(HypervisorDriver):
|
||
"""NAPALM driver for Proxmox VE nodes."""
|
||
|
||
platform = "proxmox"
|
||
|
||
def __init__(
|
||
self,
|
||
hostname: str,
|
||
username: str,
|
||
password: str,
|
||
timeout: int = 60,
|
||
optional_args: _JsonDict | None = None,
|
||
) -> None:
|
||
self.hostname = hostname
|
||
self.username = username
|
||
self.password = password
|
||
self.timeout = timeout
|
||
self.optional_args: _JsonDict = optional_args or {}
|
||
|
||
self._port: int = self.optional_args.get("port", 8006)
|
||
self._verify_ssl: bool = self.optional_args.get(
|
||
"verify_ssl", self.optional_args.get("ssl_verify", self.optional_args.get("verify", True))
|
||
)
|
||
self._realm: str = self.optional_args.get("realm", "pam")
|
||
self._token_name: str | None = self.optional_args.get("token_name")
|
||
self._token_value: str | None = self.optional_args.get("token_value")
|
||
self._node: str | None = self.optional_args.get("node")
|
||
|
||
self._api: ProxmoxAPI | None = None
|
||
self._node_name: str = ""
|
||
self._ssh_client: "paramiko.SSHClient | None" = None
|
||
|
||
# Candidate config (merge/replace)
|
||
self._candidate_config: str = ""
|
||
self._running_config: str = ""
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Connection management
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def open(self) -> None:
|
||
"""Open the connection to the Proxmox API."""
|
||
try:
|
||
kwargs: _JsonDict = {
|
||
"host": self.hostname,
|
||
"port": self._port,
|
||
"verify_ssl": self._verify_ssl,
|
||
"timeout": self.timeout,
|
||
}
|
||
|
||
if self._token_name and self._token_value:
|
||
# Token-based auth: token_name format is "user@realm!tokenid"
|
||
# proxmoxer needs user="user@realm", token_name="tokenid", token_value="..."
|
||
if "!" in self._token_name:
|
||
user_part, token_id = self._token_name.split("!", 1)
|
||
else:
|
||
user_part = self.username or "root@pam"
|
||
token_id = self._token_name
|
||
kwargs["user"] = user_part
|
||
kwargs["token_name"] = token_id
|
||
kwargs["token_value"] = self._token_value
|
||
else:
|
||
kwargs["user"] = f"{self.username}@{self._realm}"
|
||
kwargs["password"] = self.password
|
||
|
||
self._api = ProxmoxAPI(**kwargs)
|
||
|
||
# Determine which node we're talking to
|
||
self._node_name = self._resolve_node()
|
||
except Exception as exc:
|
||
raise ConnectionException(
|
||
f"Cannot connect to Proxmox at {self.hostname}:{self._port} — {exc}"
|
||
) from exc
|
||
|
||
# Establish a persistent SSH connection for shell commands (used by
|
||
# _exec_ssh_command as a fast path, avoiding per-call connect overhead).
|
||
if self.username and self.password:
|
||
try:
|
||
import paramiko # noqa: PLC0415
|
||
ssh_port: int = self.optional_args.get("ssh_port", 22)
|
||
client = paramiko.SSHClient()
|
||
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
|
||
client.connect(
|
||
hostname=self.hostname,
|
||
port=ssh_port,
|
||
username=self.username,
|
||
password=self.password,
|
||
timeout=self.timeout,
|
||
look_for_keys=False,
|
||
allow_agent=False,
|
||
)
|
||
self._ssh_client = client
|
||
except Exception:
|
||
self._ssh_client = None
|
||
|
||
def _resolve_node(self) -> str:
|
||
"""Return the Proxmox node name for this host."""
|
||
if self._node:
|
||
return self._node
|
||
# Try to match hostname against listed nodes
|
||
try:
|
||
nodes = self._api.nodes.get() # type: ignore[union-attr]
|
||
except Exception:
|
||
nodes = []
|
||
short_host = self.hostname.split(".")[0].lower()
|
||
for node_entry in nodes:
|
||
name = node_entry.get("node", "")
|
||
if name.lower() == short_host or name.lower() == self.hostname.lower():
|
||
return name
|
||
# Fall back to first online node
|
||
for node_entry in nodes:
|
||
if node_entry.get("status") == "online":
|
||
return node_entry["node"]
|
||
# Last resort: use the shortened hostname
|
||
return short_host
|
||
|
||
def close(self) -> None:
|
||
"""Close the session (Proxmox REST is stateless; nothing to tear down)."""
|
||
self._api = None
|
||
if self._ssh_client is not None:
|
||
try:
|
||
self._ssh_client.close()
|
||
except Exception:
|
||
pass
|
||
self._ssh_client = None
|
||
|
||
def is_alive(self) -> _JsonDict:
|
||
"""Return connection state."""
|
||
try:
|
||
self._api.version.get() # type: ignore[union-attr]
|
||
return {"is_alive": True}
|
||
except Exception:
|
||
return {"is_alive": False}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Internal helpers
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def _node_api(self):
|
||
"""Return the proxmoxer sub-resource for the target node."""
|
||
return self._api.nodes(self._node_name) # type: ignore[union-attr]
|
||
|
||
def _get_node_network(self) -> list[_JsonDict]:
|
||
"""Return the list of network interfaces from the Proxmox node API."""
|
||
try:
|
||
return self._node_api().network.get() or []
|
||
except ResourceException:
|
||
return []
|
||
|
||
def _exec_ssh_command(self, command: str) -> str:
|
||
"""Execute a shell command on the Proxmox node.
|
||
|
||
First tries the Proxmox API execute endpoint. If that fails (e.g.
|
||
because token-auth is blocked on that endpoint), falls back to SSH
|
||
using the driver's username / password credentials.
|
||
Falls back to an empty string if both methods are unavailable.
|
||
"""
|
||
# 1. Try Proxmox API execute endpoint
|
||
try:
|
||
result = self._node_api().execute.post(command=command)
|
||
data = result.get("data", "")
|
||
if isinstance(data, str):
|
||
return data
|
||
except Exception:
|
||
pass
|
||
|
||
# 2. Fall back to SSH using explicit SSH credentials from optional_args,
|
||
# or the driver's own username/password as a last resort.
|
||
ssh_user = self.optional_args.get("ssh_username") or self.username
|
||
ssh_pass = self.optional_args.get("ssh_password") or self.password
|
||
ssh_key_str: str | None = self.optional_args.get("ssh_private_key_str")
|
||
|
||
if not ssh_user or (not ssh_pass and not ssh_key_str):
|
||
return ""
|
||
try:
|
||
import paramiko # noqa: PLC0415
|
||
import tempfile, os as _os # noqa: PLC0415
|
||
client = self._ssh_client
|
||
if client is None or not client.get_transport() or not client.get_transport().is_active():
|
||
ssh_port: int = self.optional_args.get("ssh_port", 22)
|
||
client = paramiko.SSHClient()
|
||
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
|
||
connect_kwargs: dict = dict(
|
||
hostname=self.hostname,
|
||
port=ssh_port,
|
||
username=ssh_user,
|
||
timeout=self.timeout,
|
||
look_for_keys=False,
|
||
allow_agent=False,
|
||
)
|
||
_tmp_key = None
|
||
if ssh_key_str:
|
||
_tmp_key = tempfile.NamedTemporaryFile(mode="w", suffix=".pem", delete=False)
|
||
_tmp_key.write(ssh_key_str)
|
||
_tmp_key.flush()
|
||
_tmp_key.close()
|
||
connect_kwargs["key_filename"] = _tmp_key.name
|
||
else:
|
||
connect_kwargs["password"] = ssh_pass
|
||
try:
|
||
client.connect(**connect_kwargs)
|
||
finally:
|
||
if _tmp_key:
|
||
try:
|
||
_os.unlink(_tmp_key.name)
|
||
except OSError:
|
||
pass
|
||
self._ssh_client = client
|
||
_, stdout, _ = client.exec_command(command, timeout=self.timeout)
|
||
output = stdout.read().decode("utf-8", errors="replace")
|
||
return output
|
||
except Exception:
|
||
self._ssh_client = None
|
||
return ""
|
||
|
||
def _get_version_info(self) -> _JsonDict:
|
||
"""Return Proxmox version dict."""
|
||
try:
|
||
return self._api.version.get() or {} # type: ignore[union-attr]
|
||
except Exception:
|
||
return {}
|
||
|
||
def _get_node_status(self) -> _JsonDict:
|
||
"""Return node status dict."""
|
||
try:
|
||
return self._node_api().status.get() or {}
|
||
except Exception:
|
||
return {}
|
||
|
||
def _get_node_subscription(self) -> _JsonDict:
|
||
try:
|
||
return self._node_api().subscription.get() or {}
|
||
except Exception:
|
||
return {}
|
||
|
||
def _get_sdn_zones(self) -> list[_JsonDict]:
|
||
try:
|
||
return self._api.cluster.sdn.zones.get() or [] # type: ignore[union-attr]
|
||
except Exception:
|
||
return []
|
||
|
||
def _get_sdn_vnets(self) -> list[_JsonDict]:
|
||
try:
|
||
return self._api.cluster.sdn.vnets.get() or [] # type: ignore[union-attr]
|
||
except Exception:
|
||
return []
|
||
|
||
def _get_sdn_subnets(self, vnet: str) -> list[_JsonDict]:
|
||
try:
|
||
return self._api.cluster.sdn.vnets(vnet).subnets.get() or []
|
||
except Exception:
|
||
return []
|
||
|
||
def _get_node_dns(self) -> _JsonDict:
|
||
try:
|
||
return self._node_api().dns.get() or {}
|
||
except Exception:
|
||
return {}
|
||
|
||
def _get_node_time(self) -> _JsonDict:
|
||
try:
|
||
return self._node_api().time.get() or {}
|
||
except Exception:
|
||
return {}
|
||
|
||
def _get_node_ntp(self) -> _JsonDict:
|
||
try:
|
||
return self._node_api().ntp.get() or {}
|
||
except Exception:
|
||
return {}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_facts
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_facts(self) -> _JsonDict:
|
||
"""Return basic facts about the Proxmox node."""
|
||
status = self._get_node_status()
|
||
version = self._get_version_info()
|
||
network = self._get_node_network()
|
||
dns = self._get_node_dns()
|
||
|
||
uptime = float(status.get("uptime", 0))
|
||
model = status.get("model", "")
|
||
# dns["search"] is the DNS search domain (e.g. "home.example.com"), not the hostname
|
||
dns_search = dns.get("search", "")
|
||
hostname = self._node_name
|
||
fqdn = f"{self._node_name}.{dns_search}" if dns_search else self.hostname
|
||
|
||
# Build interface list
|
||
iface_list = sorted(
|
||
iface["iface"] for iface in network if iface.get("iface")
|
||
)
|
||
|
||
# PVE version looks like "8.2.4"
|
||
pve_version = version.get("version", "")
|
||
release = version.get("release", "")
|
||
os_version = f"Proxmox VE {pve_version}" if pve_version else f"Proxmox VE {release}"
|
||
|
||
return {
|
||
"uptime": uptime,
|
||
"vendor": "Proxmox Server Solutions GmbH",
|
||
"model": model or "Proxmox VE Node",
|
||
"hostname": self._node_name,
|
||
"fqdn": fqdn or self.hostname,
|
||
"os_version": os_version,
|
||
"serial_number": "",
|
||
"interface_list": iface_list,
|
||
}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_interfaces
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_interfaces(self) -> dict[str, _JsonDict]:
|
||
"""Return a dict of interfaces keyed by interface name."""
|
||
result: dict[str, _JsonDict] = {}
|
||
for iface in self._get_node_network():
|
||
name = iface.get("iface", "")
|
||
if not name:
|
||
continue
|
||
|
||
# Proxmox marks bridge ports as not independent — include all
|
||
active = iface.get("active", 0)
|
||
autostart = iface.get("autostart", 0)
|
||
|
||
result[name] = {
|
||
"is_up": bool(active),
|
||
"is_enabled": bool(autostart) or bool(active),
|
||
"description": iface.get("comments", "").strip(),
|
||
"last_flapped": -1.0,
|
||
"speed": utils.speed_mbps(iface),
|
||
"mtu": int(iface.get("mtu") or 1500),
|
||
"mac_address": utils.normalize_mac(iface.get("hwaddr", "")),
|
||
}
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_interfaces_ip
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_interfaces_ip(self) -> dict[str, _JsonDict]:
|
||
"""Return IP addresses per interface."""
|
||
result: dict[str, _JsonDict] = {}
|
||
for iface in self._get_node_network():
|
||
name = iface.get("iface", "")
|
||
if not name:
|
||
continue
|
||
addrs = utils.addresses_from_node_network(iface)
|
||
if addrs:
|
||
result[name] = addrs
|
||
|
||
# Overlay SDN VNet addresses (subnets with gateways)
|
||
for vnet in self._get_sdn_vnets():
|
||
vnet_id = vnet.get("vnet", "")
|
||
if not vnet_id:
|
||
continue
|
||
for subnet in self._get_sdn_subnets(vnet_id):
|
||
cidr = subnet.get("cidr", "")
|
||
gateway = subnet.get("gateway", "")
|
||
if gateway and cidr:
|
||
ip, plen = utils.parse_cidr(cidr)
|
||
if "." in gateway:
|
||
result.setdefault(vnet_id, {}).setdefault("ipv4", {})[gateway] = {
|
||
"prefix_length": plen
|
||
}
|
||
else:
|
||
result.setdefault(vnet_id, {}).setdefault("ipv6", {})[gateway] = {
|
||
"prefix_length": plen
|
||
}
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_interfaces_counters
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_interfaces_counters(self) -> dict[str, _JsonDict]:
|
||
"""Return per-interface traffic counters."""
|
||
result: dict[str, _JsonDict] = {}
|
||
try:
|
||
rrd_data = self._node_api().netstat.get() or []
|
||
except Exception:
|
||
rrd_data = []
|
||
|
||
# Proxmox /nodes/{node}/netstat returns a list of time-series points.
|
||
# Take the most recent (last) entry for each interface.
|
||
latest: dict[str, _JsonDict] = {}
|
||
for entry in rrd_data:
|
||
iface = entry.get("dev", "")
|
||
if iface:
|
||
latest[iface] = entry
|
||
|
||
for iface, data in latest.items():
|
||
result[iface] = {
|
||
"tx_errors": int(data.get("tx_errs", 0) or 0),
|
||
"rx_errors": int(data.get("rx_errs", 0) or 0),
|
||
"tx_discards": int(data.get("tx_drop", 0) or 0),
|
||
"rx_discards": int(data.get("rx_drop", 0) or 0),
|
||
"tx_octets": int(data.get("tx_bytes", 0) or 0),
|
||
"rx_octets": int(data.get("rx_bytes", 0) or 0),
|
||
"tx_unicast_packets": int(data.get("tx_packets", 0) or 0),
|
||
"rx_unicast_packets": int(data.get("rx_packets", 0) or 0),
|
||
"tx_multicast_packets": 0,
|
||
"rx_multicast_packets": 0,
|
||
"tx_broadcast_packets": 0,
|
||
"rx_broadcast_packets": 0,
|
||
}
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_environment
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_environment(self) -> _JsonDict:
|
||
"""Return environment status (CPU, memory, temperature)."""
|
||
status = self._get_node_status()
|
||
env: _JsonDict = {
|
||
"fans": {},
|
||
"temperature": {},
|
||
"power": {},
|
||
"cpu": {},
|
||
"memory": {"available_ram": 0, "used_ram": 0},
|
||
}
|
||
|
||
# CPU
|
||
cpu_usage = status.get("cpu", 0.0)
|
||
env["cpu"]["0"] = {"%usage": round(float(cpu_usage) * 100, 2)}
|
||
|
||
# Memory (Proxmox reports in bytes)
|
||
mem = status.get("memory", {})
|
||
total = int(mem.get("total", 0) or 0)
|
||
used = int(mem.get("used", 0) or 0)
|
||
env["memory"]["available_ram"] = total
|
||
env["memory"]["used_ram"] = used
|
||
|
||
# Temperature (from node sensors if available)
|
||
try:
|
||
sensors = self._node_api().hardware.sensors.get() or []
|
||
except Exception:
|
||
sensors = []
|
||
for sensor in sensors:
|
||
name = sensor.get("name", "unknown")
|
||
value = sensor.get("value", None)
|
||
if value is not None:
|
||
try:
|
||
temp_c = float(value)
|
||
env["temperature"][name] = {
|
||
"temperature": temp_c,
|
||
"is_alert": temp_c >= 80.0,
|
||
"is_critical": temp_c >= 95.0,
|
||
}
|
||
except (TypeError, ValueError):
|
||
pass
|
||
|
||
return env
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_arp_table
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_arp_table(self, vrf: str = "") -> list[_JsonDict]:
|
||
"""Return ARP table.
|
||
|
||
Proxmox does not expose ARP via the REST API directly. We attempt
|
||
to read it via the node's ``/proc/net/arp`` through the Proxmox
|
||
exec endpoint. If that is unavailable, an empty list is returned.
|
||
"""
|
||
raw = self._exec_ssh_command("cat /proc/net/arp")
|
||
if not raw:
|
||
return []
|
||
|
||
entries = []
|
||
# /proc/net/arp format:
|
||
# IP address HW type Flags HW address Mask Device
|
||
# 192.168.1.1 0x1 0x2 aa:bb:cc:dd:ee:ff * eth0
|
||
for line in raw.splitlines():
|
||
line = line.strip()
|
||
if not line or line.startswith("IP"):
|
||
continue
|
||
parts = line.split()
|
||
if len(parts) < 6:
|
||
continue
|
||
ip_addr, _, flags, mac, _, iface = (
|
||
parts[0], parts[1], parts[2], parts[3], parts[4], parts[5]
|
||
)
|
||
if mac in ("00:00:00:00:00:00", ""):
|
||
continue
|
||
if vrf and iface != vrf:
|
||
continue
|
||
entries.append(
|
||
{
|
||
"interface": iface,
|
||
"mac": utils.normalize_mac(mac),
|
||
"ip": ip_addr,
|
||
"age": -1.0,
|
||
}
|
||
)
|
||
return entries
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_mac_address_table
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_mac_address_table(self) -> list[_JsonDict]:
|
||
"""Return MAC address table from Linux bridges and OVS bridges."""
|
||
result: list[_JsonDict] = []
|
||
network = self._get_node_network()
|
||
|
||
# Linux bridges
|
||
linux_bridges = [
|
||
iface["iface"]
|
||
for iface in network
|
||
if iface.get("type") in ("bridge",) and iface.get("iface")
|
||
]
|
||
for bridge in linux_bridges:
|
||
raw = self._exec_ssh_command(
|
||
f"bridge fdb show br {bridge} 2>/dev/null || true"
|
||
)
|
||
for line in raw.splitlines():
|
||
parts = line.split()
|
||
if len(parts) < 3:
|
||
continue
|
||
mac_str = parts[0]
|
||
if not re.match(r"([0-9a-f]{2}:){5}[0-9a-f]{2}", mac_str):
|
||
continue
|
||
# "dev <port>" "vlan <vid>"
|
||
dev = ""
|
||
vlan_id = 1
|
||
for i, tok in enumerate(parts):
|
||
if tok == "dev" and i + 1 < len(parts):
|
||
dev = parts[i + 1]
|
||
if tok == "vlan" and i + 1 < len(parts):
|
||
try:
|
||
vlan_id = int(parts[i + 1])
|
||
except ValueError:
|
||
pass
|
||
result.append(
|
||
{
|
||
"mac": utils.normalize_mac(mac_str),
|
||
"interface": dev or bridge,
|
||
"vlan": vlan_id,
|
||
"static": "permanent" in line,
|
||
"active": True,
|
||
"moves": 0,
|
||
"last_move": 0.0,
|
||
}
|
||
)
|
||
|
||
# OVS bridges
|
||
ovs_bridges = [
|
||
iface["iface"]
|
||
for iface in network
|
||
if iface.get("type") in ("OVSBridge",) and iface.get("iface")
|
||
]
|
||
for bridge in ovs_bridges:
|
||
raw = self._exec_ssh_command(
|
||
f"ovs-appctl fdb/show {bridge} 2>/dev/null || true"
|
||
)
|
||
# Format: port VLAN MAC AGE
|
||
for line in raw.splitlines():
|
||
parts = line.split()
|
||
if len(parts) < 4:
|
||
continue
|
||
try:
|
||
_port = int(parts[0])
|
||
vlan_id = int(parts[1])
|
||
mac_str = parts[2]
|
||
except (ValueError, IndexError):
|
||
continue
|
||
result.append(
|
||
{
|
||
"mac": utils.normalize_mac(mac_str),
|
||
"interface": bridge,
|
||
"vlan": vlan_id,
|
||
"static": False,
|
||
"active": True,
|
||
"moves": 0,
|
||
"last_move": 0.0,
|
||
}
|
||
)
|
||
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_vlans (SDN VNets + classic Linux bridge VLANs)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_vlans(self) -> dict[str, _JsonDict]:
|
||
"""Return VLAN table.
|
||
|
||
For OVS+SDN nodes: reads SDN VNets for VLAN IDs/names, then maps
|
||
OVSIntPort (access ports with ovs_tag) → untagged membership, and
|
||
OVSPort / OVSBridge (trunk ports) → tagged membership.
|
||
|
||
Falls back to ``bridge vlan show`` for classic Linux-bridge nodes.
|
||
"""
|
||
result: dict[str, _JsonDict] = {}
|
||
node_network = self._get_node_network()
|
||
|
||
# --- SDN vnets → VLAN IDs and initial entries ---
|
||
for vnet in self._get_sdn_vnets():
|
||
tag = vnet.get("tag")
|
||
vnet_id = vnet.get("vnet", "")
|
||
if tag is None:
|
||
continue
|
||
try:
|
||
tag_int = int(tag)
|
||
except (ValueError, TypeError):
|
||
continue
|
||
result[str(tag_int)] = {
|
||
"name": vnet_id,
|
||
"tagged": [],
|
||
"untagged": [],
|
||
}
|
||
|
||
# --- OVS port membership ---
|
||
trunk_ports: list[str] = []
|
||
access_by_vlan: dict[str, list[str]] = {}
|
||
|
||
for iface in node_network:
|
||
ovs_type = iface.get("ovs_type", "")
|
||
iface_name = iface.get("iface", "")
|
||
if not iface_name:
|
||
continue
|
||
if ovs_type == "OVSIntPort":
|
||
ovs_tag = iface.get("ovs_tag")
|
||
if ovs_tag is not None:
|
||
vid = str(int(ovs_tag))
|
||
access_by_vlan.setdefault(vid, []).append(iface_name)
|
||
elif ovs_type in ("OVSPort", "OVSBridge"):
|
||
# Trunk: carries all VLANs tagged
|
||
trunk_ports.append(iface_name)
|
||
|
||
if trunk_ports or access_by_vlan:
|
||
# OVS topology detected — assign tagged/untagged per VLAN
|
||
for vid, vlan_entry in result.items():
|
||
vlan_entry["tagged"] = list(trunk_ports)
|
||
vlan_entry["untagged"] = list(access_by_vlan.get(vid, []))
|
||
return result
|
||
|
||
# --- Linux bridge fallback (bridge vlan show) ---
|
||
# Convert result entries to use tagged/untagged keys
|
||
for entry in result.values():
|
||
entry.setdefault("tagged", [])
|
||
entry.setdefault("untagged", [])
|
||
|
||
raw = self._exec_ssh_command("bridge vlan show 2>/dev/null || true")
|
||
current_iface = ""
|
||
for line in raw.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
m = re.match(r"^(\S+)\s+(\d+)", line)
|
||
if m:
|
||
current_iface = m.group(1)
|
||
vid = m.group(2)
|
||
else:
|
||
m2 = re.match(r"^\s*(\d+)", line)
|
||
if m2:
|
||
vid = m2.group(1)
|
||
else:
|
||
continue
|
||
if current_iface and vid:
|
||
entry = result.setdefault(vid, {"name": "", "tagged": [], "untagged": []})
|
||
if "Untagged" in line or "PVID" in line:
|
||
if current_iface not in entry["untagged"]:
|
||
entry["untagged"].append(current_iface)
|
||
else:
|
||
if current_iface not in entry["tagged"]:
|
||
entry["tagged"].append(current_iface)
|
||
|
||
# Only return VLANs that are actually assigned to at least one interface.
|
||
# Linux bridge vlan show reports all 4094 possible VIDs per port —
|
||
# filtering here avoids polluting the VLAN table with phantom entries.
|
||
return {
|
||
vid: entry for vid, entry in result.items()
|
||
if entry.get("tagged") or entry.get("untagged")
|
||
}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_network_instances (SDN Zones as VRF-like instances)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_network_instances(self, name: str = "") -> dict[str, _JsonDict]:
|
||
"""Return SDN zones as network instances."""
|
||
result: dict[str, _JsonDict] = {}
|
||
|
||
# Always include the default instance
|
||
result["default"] = {
|
||
"name": "default",
|
||
"type": "DEFAULT_INSTANCE",
|
||
"state": {"route_distinguisher": None},
|
||
"interfaces": {"interface": {}},
|
||
}
|
||
|
||
# Attach non-SDN interfaces to the default instance
|
||
for iface in self._get_node_network():
|
||
iface_name = iface.get("iface", "")
|
||
if iface_name:
|
||
result["default"]["interfaces"]["interface"][iface_name] = {}
|
||
|
||
# SDN Zones
|
||
for zone in self._get_sdn_zones():
|
||
zone_id = zone.get("zone", zone.get("name", ""))
|
||
if not zone_id:
|
||
continue
|
||
if name and zone_id != name:
|
||
continue
|
||
instance = utils.sdn_zone_to_network_instance(zone)
|
||
# Attach VNets that belong to this zone
|
||
for vnet in self._get_sdn_vnets():
|
||
if vnet.get("zone") == zone_id:
|
||
vnet_id = vnet.get("vnet", "")
|
||
if vnet_id:
|
||
instance["interfaces"]["interface"][vnet_id] = {}
|
||
result[zone_id] = instance
|
||
|
||
if name:
|
||
return {k: v for k, v in result.items() if k == name}
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_ntp_servers / get_ntp_stats
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_ntp_servers(self) -> dict[str, _JsonDict]:
|
||
"""Return configured NTP servers."""
|
||
ntp = self._get_node_ntp()
|
||
servers: dict[str, _JsonDict] = {}
|
||
# Proxmox reports a comma-separated or space-separated server list
|
||
raw = ntp.get("server", "") or ntp.get("servers", "")
|
||
for srv in re.split(r"[\s,]+", raw):
|
||
srv = srv.strip()
|
||
if srv:
|
||
servers[srv] = {}
|
||
return servers
|
||
|
||
def get_ntp_stats(self) -> list[_JsonDict]:
|
||
"""Return NTP synchronisation statistics from chronyc/ntpq output."""
|
||
raw = self._exec_ssh_command(
|
||
"chronyc -n tracking 2>/dev/null || ntpq -pn 2>/dev/null || true"
|
||
)
|
||
stats: list[_JsonDict] = []
|
||
for line in raw.splitlines():
|
||
line = line.strip()
|
||
# ntpq -pn format: *remote refid st t when poll reach delay offset jitter
|
||
m = re.match(
|
||
r"^([\*\+\-\s])([\d.]+)\s+([\d.]+)\s+(\d+)\s+\S+\s+(\S+)\s+(\d+)\s+(\d+)\s+([\d.]+)\s+([-\d.]+)\s+([\d.]+)",
|
||
line,
|
||
)
|
||
if m:
|
||
synced = m.group(1).strip() == "*"
|
||
stats.append(
|
||
{
|
||
"remote": m.group(2),
|
||
"referenceid": m.group(3),
|
||
"synchronized": synced,
|
||
"stratum": int(m.group(4)),
|
||
"type": "",
|
||
"when": m.group(5),
|
||
"hostpoll": int(m.group(6)),
|
||
"reachability": int(m.group(7)),
|
||
"delay": float(m.group(8)),
|
||
"offset": float(m.group(9)),
|
||
"jitter": float(m.group(10)),
|
||
}
|
||
)
|
||
return stats
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_snmp_information
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_snmp_information(self) -> _JsonDict:
|
||
"""Return SNMP information.
|
||
|
||
Proxmox does not expose SNMP configuration via the REST API.
|
||
We read /etc/snmp/snmpd.conf via exec if available.
|
||
"""
|
||
raw = self._exec_ssh_command(
|
||
"cat /etc/snmp/snmpd.conf 2>/dev/null || true"
|
||
)
|
||
communities: dict[str, _JsonDict] = {}
|
||
location = ""
|
||
contact = ""
|
||
for line in raw.splitlines():
|
||
line = line.strip()
|
||
if line.startswith("#") or not line:
|
||
continue
|
||
# rocommunity <community> [source]
|
||
m = re.match(r"^(ro|rw)community\s+(\S+)", line)
|
||
if m:
|
||
mode = "ro" if m.group(1) == "ro" else "rw"
|
||
community = m.group(2)
|
||
communities[community] = {"acl": "N/A", "mode": mode}
|
||
m_loc = re.match(r"^sysLocation\s+(.+)", line)
|
||
if m_loc:
|
||
location = m_loc.group(1).strip()
|
||
m_con = re.match(r"^sysContact\s+(.+)", line)
|
||
if m_con:
|
||
contact = m_con.group(1).strip()
|
||
|
||
return {
|
||
"chassis_id": self._node_name,
|
||
"community": communities,
|
||
"contact": contact,
|
||
"location": location,
|
||
}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_users
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_users(self) -> dict[str, _JsonDict]:
|
||
"""Return users configured on the Proxmox node.
|
||
|
||
Reads from both the Proxmox access/users API and local /etc/passwd.
|
||
"""
|
||
result: dict[str, _JsonDict] = {}
|
||
try:
|
||
pve_users = self._api.access.users.get() or [] # type: ignore[union-attr]
|
||
except Exception:
|
||
pve_users = []
|
||
|
||
for user in pve_users:
|
||
uid = user.get("userid", "")
|
||
if not uid:
|
||
continue
|
||
# Proxmox roles: Administrator → 15, otherwise 1
|
||
groups = user.get("groups", "") or ""
|
||
level = 1
|
||
try:
|
||
roles = self._api.access.users(uid).get() or {} # type: ignore[union-attr]
|
||
if "Administrator" in str(roles):
|
||
level = 15
|
||
except Exception:
|
||
pass
|
||
result[uid] = {
|
||
"level": level,
|
||
"password": "",
|
||
"sshkeys": [],
|
||
}
|
||
|
||
# Merge local OS users from /etc/passwd
|
||
raw = self._exec_ssh_command("getent passwd 2>/dev/null || cat /etc/passwd")
|
||
for line in raw.splitlines():
|
||
parts = line.split(":")
|
||
if len(parts) < 7:
|
||
continue
|
||
uname, _, uid_str, *_ = parts
|
||
try:
|
||
uid_int = int(uid_str)
|
||
except ValueError:
|
||
continue
|
||
if uname not in result and uid_int < 1000 or uid_int == 0:
|
||
result[uname] = {
|
||
"level": 15 if uid_int == 0 else 0,
|
||
"password": "",
|
||
"sshkeys": [],
|
||
}
|
||
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_config
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_config(
|
||
self,
|
||
retrieve: str = "all",
|
||
full: bool = False,
|
||
sanitized: bool = False,
|
||
format: str = "text",
|
||
) -> _JsonDict:
|
||
"""Return the node network configuration.
|
||
|
||
``running`` config is the contents of ``/etc/network/interfaces``
|
||
(and the SDN config directory). ``startup`` is identical (PVE
|
||
applies on boot). ``candidate`` is what was loaded via
|
||
``load_merge_candidate`` / ``load_replace_candidate`` but not yet
|
||
committed.
|
||
"""
|
||
configs: _JsonDict = {"running": "", "candidate": "", "startup": ""}
|
||
|
||
if retrieve in ("running", "all", "startup"):
|
||
raw = self._exec_ssh_command(
|
||
"cat /etc/network/interfaces 2>/dev/null || true"
|
||
)
|
||
# Append SDN config if available
|
||
sdn_raw = self._exec_ssh_command(
|
||
"cat /etc/pve/sdn/vnets.cfg 2>/dev/null || true"
|
||
)
|
||
running = raw
|
||
if sdn_raw:
|
||
running += "\n# === SDN VNets ===\n" + sdn_raw
|
||
if sanitized:
|
||
running = re.sub(r"password\s+\S+", "password ****", running)
|
||
configs["running"] = running
|
||
self._running_config = running
|
||
if retrieve == "startup":
|
||
configs["startup"] = running
|
||
elif retrieve == "all":
|
||
configs["startup"] = running
|
||
|
||
if retrieve in ("candidate", "all"):
|
||
configs["candidate"] = self._candidate_config
|
||
|
||
return configs
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Config management (load / compare / commit / discard / rollback)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def load_merge_candidate(
|
||
self,
|
||
filename: str | None = None,
|
||
config: str | None = None,
|
||
) -> None:
|
||
"""Load a candidate configuration (merge mode)."""
|
||
if filename:
|
||
with open(filename) as fh:
|
||
config = fh.read()
|
||
if config is None:
|
||
raise ValueError("Either filename or config must be provided")
|
||
# In merge mode we append / overlay
|
||
self._candidate_config = config
|
||
|
||
def load_replace_candidate(
|
||
self,
|
||
filename: str | None = None,
|
||
config: str | None = None,
|
||
) -> None:
|
||
"""Load a candidate configuration (replace mode)."""
|
||
if filename:
|
||
with open(filename) as fh:
|
||
config = fh.read()
|
||
if config is None:
|
||
raise ValueError("Either filename or config must be provided")
|
||
self._candidate_config = config
|
||
|
||
def compare_config(self) -> str:
|
||
"""Return a unified diff between running and candidate config."""
|
||
import difflib
|
||
|
||
if not self._running_config:
|
||
self.get_config(retrieve="running")
|
||
running_lines = self._running_config.splitlines(keepends=True)
|
||
candidate_lines = self._candidate_config.splitlines(keepends=True)
|
||
diff = difflib.unified_diff(
|
||
running_lines,
|
||
candidate_lines,
|
||
fromfile="running",
|
||
tofile="candidate",
|
||
)
|
||
return "".join(diff)
|
||
|
||
def commit_config(self, message: str = "", revert_in: int | None = None) -> None:
|
||
"""Commit the candidate configuration to the Proxmox node.
|
||
|
||
This writes the candidate config to ``/etc/network/interfaces``
|
||
via the Proxmox node/network PUT API (which applies it live).
|
||
|
||
.. note::
|
||
Full programmatic apply requires the Proxmox API to accept raw
|
||
interface configs. This implementation uses ``pvesh`` via exec
|
||
which requires the node exec endpoint to be available.
|
||
"""
|
||
if not self._candidate_config:
|
||
return
|
||
# Write via exec endpoint
|
||
escaped = self._candidate_config.replace("'", "'\\''")
|
||
self._exec_ssh_command(
|
||
f"printf '%s' '{escaped}' > /etc/network/interfaces && "
|
||
"ifreload -a 2>&1 || ifup -a 2>&1 || true"
|
||
)
|
||
self._running_config = self._candidate_config
|
||
self._candidate_config = ""
|
||
|
||
def discard_config(self) -> None:
|
||
"""Discard the candidate configuration."""
|
||
self._candidate_config = ""
|
||
|
||
def rollback(self) -> None:
|
||
"""Revert to the stored running configuration."""
|
||
if self._running_config:
|
||
self._candidate_config = self._running_config
|
||
self.commit_config()
|
||
self._candidate_config = ""
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_lldp_neighbors (via lldpcli if installed)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_lldp_neighbors(self) -> dict[str, list[_JsonDict]]:
|
||
"""Return LLDP neighbours (requires lldpd on the Proxmox node)."""
|
||
result: dict[str, list[_JsonDict]] = {}
|
||
raw = self._exec_ssh_command(
|
||
"lldpcli show neighbors summary 2>/dev/null || true"
|
||
)
|
||
current_iface = ""
|
||
for line in raw.splitlines():
|
||
m_iface = re.match(r"^\s*Interface:\s+(\S+?),?\s", line)
|
||
if m_iface:
|
||
current_iface = m_iface.group(1)
|
||
result.setdefault(current_iface, [])
|
||
continue
|
||
m_sys = re.match(r"^\s*SysName:\s+(.+)", line)
|
||
m_port = re.match(r"^\s*PortID:\s+\S+\s+(.+)", line)
|
||
if m_sys and current_iface:
|
||
hostname = m_sys.group(1).strip()
|
||
if result[current_iface]:
|
||
result[current_iface][-1]["hostname"] = hostname
|
||
else:
|
||
result[current_iface].append({"hostname": hostname, "port": ""})
|
||
if m_port and current_iface and result[current_iface]:
|
||
result[current_iface][-1]["port"] = m_port.group(1).strip()
|
||
return result
|
||
|
||
def get_lldp_neighbors_detail(self, interface: str = "") -> dict[str, list[_JsonDict]]:
|
||
"""Return detailed LLDP neighbour information."""
|
||
result: dict[str, list[_JsonDict]] = {}
|
||
raw = self._exec_ssh_command(
|
||
"lldpcli show neighbors details 2>/dev/null || true"
|
||
)
|
||
current_iface = ""
|
||
current_entry: _JsonDict = {}
|
||
|
||
def _flush():
|
||
if current_iface and current_entry:
|
||
result.setdefault(current_iface, []).append(current_entry.copy())
|
||
|
||
for line in raw.splitlines():
|
||
m_iface = re.match(r"^\s*Interface:\s+(\S+?),?\s", line)
|
||
if m_iface:
|
||
_flush()
|
||
current_iface = m_iface.group(1)
|
||
if interface and current_iface != interface:
|
||
current_iface = ""
|
||
current_entry = {
|
||
"parent_interface": "",
|
||
"remote_chassis_id": "",
|
||
"remote_system_name": "",
|
||
"remote_port": "",
|
||
"remote_port_description": "",
|
||
"remote_system_description": "",
|
||
"remote_system_capab": [],
|
||
"remote_system_enable_capab": [],
|
||
}
|
||
continue
|
||
if not current_iface:
|
||
continue
|
||
for key, pattern in (
|
||
("remote_chassis_id", r"ChassisID:\s+\S+\s+(.+)"),
|
||
("remote_system_name", r"SysName:\s+(.+)"),
|
||
("remote_port", r"PortID:\s+\S+\s+(.+)"),
|
||
("remote_port_description", r"PortDescr:\s+(.+)"),
|
||
("remote_system_description", r"SysDescr:\s+(.+)"),
|
||
):
|
||
m = re.match(rf"^\s*{pattern}", line)
|
||
if m:
|
||
current_entry[key] = m.group(1).strip()
|
||
|
||
m_cap = re.match(r"^\s*Capability:\s+(\S+),\s+(\w+)", line)
|
||
if m_cap:
|
||
cap = m_cap.group(1).lower()
|
||
enabled = m_cap.group(2).lower() == "on"
|
||
current_entry["remote_system_capab"].append(cap)
|
||
if enabled:
|
||
current_entry["remote_system_enable_capab"].append(cap)
|
||
|
||
_flush()
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# get_ipv6_neighbors_table
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_ipv6_neighbors_table(self) -> list[_JsonDict]:
|
||
"""Return IPv6 NDP neighbour table."""
|
||
raw = self._exec_ssh_command("ip -6 neigh show 2>/dev/null || true")
|
||
result = []
|
||
for line in raw.splitlines():
|
||
parts = line.split()
|
||
# Format: <ip> dev <iface> lladdr <mac> <state>
|
||
if len(parts) < 5:
|
||
continue
|
||
ip6 = parts[0]
|
||
iface = parts[2] if len(parts) > 2 else ""
|
||
mac = ""
|
||
state = ""
|
||
for i, tok in enumerate(parts):
|
||
if tok == "lladdr" and i + 1 < len(parts):
|
||
mac = parts[i + 1]
|
||
if tok in ("REACHABLE", "STALE", "DELAY", "PROBE", "FAILED", "NOARP", "PERMANENT"):
|
||
state = tok
|
||
if not mac or mac == "FAILED":
|
||
continue
|
||
result.append(
|
||
{
|
||
"interface": iface,
|
||
"mac": utils.normalize_mac(mac),
|
||
"ip": ip6,
|
||
"age": -1.0,
|
||
"state": state,
|
||
}
|
||
)
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# ping (via Proxmox node/execute)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def ping(
|
||
self,
|
||
destination: str,
|
||
source: str = C.PING_SOURCE,
|
||
ttl: int = C.PING_TTL,
|
||
timeout: int = C.PING_TIMEOUT,
|
||
size: int = C.PING_SIZE,
|
||
count: int = C.PING_COUNT,
|
||
vrf: str = C.PING_VRF,
|
||
source_interface: str = C.PING_SOURCE_INTERFACE,
|
||
) -> _JsonDict:
|
||
"""Execute ping on the Proxmox node and return results."""
|
||
cmd_parts = [
|
||
f"ping -c {count}",
|
||
f"-W {timeout}",
|
||
f"-s {size}",
|
||
f"-t {ttl}",
|
||
]
|
||
if source:
|
||
cmd_parts.append(f"-I {source}")
|
||
elif source_interface:
|
||
cmd_parts.append(f"-I {source_interface}")
|
||
cmd_parts.append(destination)
|
||
cmd = " ".join(cmd_parts)
|
||
|
||
raw = self._exec_ssh_command(f"{cmd} 2>&1 || true")
|
||
if not raw:
|
||
return {"error": "Ping command not available via exec endpoint"}
|
||
|
||
# Detect common failure strings before parsing statistics
|
||
_error_patterns = (
|
||
"Name or service not known",
|
||
"Network is unreachable",
|
||
"connect: No route to host",
|
||
"unknown host",
|
||
)
|
||
for _pat in _error_patterns:
|
||
if _pat.lower() in raw.lower():
|
||
return {"error": raw.strip()}
|
||
|
||
# Parse statistics line: "5 packets transmitted, 5 received, 0% packet loss"
|
||
m_stat = re.search(
|
||
r"(\d+) packets transmitted,\s+(\d+) received.*?([\d.]+)% packet loss",
|
||
raw,
|
||
)
|
||
if not m_stat:
|
||
return {"error": raw.strip()}
|
||
|
||
sent = int(m_stat.group(1))
|
||
received = int(m_stat.group(2))
|
||
loss = sent - received
|
||
|
||
# RTT line: "rtt min/avg/max/mdev = 0.123/0.456/0.789/0.100 ms"
|
||
m_rtt = re.search(
|
||
r"rtt min/avg/max/mdev = ([\d.]+)/([\d.]+)/([\d.]+)/([\d.]+)",
|
||
raw,
|
||
)
|
||
rtt_min = float(m_rtt.group(1)) if m_rtt else 0.0
|
||
rtt_avg = float(m_rtt.group(2)) if m_rtt else 0.0
|
||
rtt_max = float(m_rtt.group(3)) if m_rtt else 0.0
|
||
rtt_std = float(m_rtt.group(4)) if m_rtt else 0.0
|
||
|
||
# Individual probe lines
|
||
probes = []
|
||
for m_probe in re.finditer(
|
||
r"icmp_seq=\d+.*?time=([\d.]+) ms.*?from ([\d.a-fA-F:]+)", raw
|
||
):
|
||
probes.append(
|
||
{"ip_address": m_probe.group(2), "rtt": float(m_probe.group(1))}
|
||
)
|
||
|
||
return {
|
||
"success": {
|
||
"probes_sent": sent,
|
||
"packet_loss": loss,
|
||
"rtt_min": rtt_min,
|
||
"rtt_max": rtt_max,
|
||
"rtt_avg": rtt_avg,
|
||
"rtt_stddev": rtt_std,
|
||
"results": probes,
|
||
}
|
||
}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# traceroute
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def traceroute(
|
||
self,
|
||
destination: str,
|
||
source: str = "",
|
||
ttl: int = 255,
|
||
timeout: int = 2,
|
||
vrf: str = "",
|
||
) -> _JsonDict:
|
||
"""Execute traceroute on the Proxmox node and return results."""
|
||
cmd_parts = [f"traceroute -m {ttl}", f"-w {timeout}", "-n"]
|
||
if source:
|
||
cmd_parts.append(f"-s {source}")
|
||
cmd_parts.append(destination)
|
||
raw = self._exec_ssh_command(" ".join(cmd_parts) + " 2>&1 || true")
|
||
|
||
if not raw:
|
||
return {"error": "traceroute not available via exec endpoint"}
|
||
|
||
hops: _JsonDict = {}
|
||
for line in raw.splitlines():
|
||
m = re.match(
|
||
r"^\s*(\d+)\s+([\d.a-fA-F:]+|\*)\s+([\d.]+|[\d.]+\s+ms|\*)",
|
||
line,
|
||
)
|
||
if not m:
|
||
continue
|
||
hop_id = int(m.group(1))
|
||
ip_addr = m.group(2)
|
||
if ip_addr == "*":
|
||
continue
|
||
# Parse RTT probes: each hop can have up to 3
|
||
rtts = re.findall(r"([\d.]+)\s+ms", line)
|
||
probes_dict = {}
|
||
for idx, rtt in enumerate(rtts, start=1):
|
||
probes_dict[idx] = {
|
||
"rtt": float(rtt),
|
||
"ip_address": ip_addr,
|
||
"host_name": ip_addr,
|
||
}
|
||
if probes_dict:
|
||
hops[hop_id] = {"probes": probes_dict}
|
||
|
||
if not hops:
|
||
return {"error": raw.strip()}
|
||
return {"success": hops}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# cli
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def cli(self, commands: list[str], encoding: str = "text") -> dict[str, str]:
|
||
"""Execute arbitrary commands on the node via the exec endpoint."""
|
||
output = {}
|
||
for cmd in commands:
|
||
output[cmd] = self._exec_ssh_command(cmd)
|
||
return output
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Unsupported / not-applicable methods
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_bgp_config(self, group: str = "", neighbor: str = "") -> _JsonDict:
|
||
raise NotImplementedError("BGP configuration is not managed via Proxmox API")
|
||
|
||
def get_bgp_neighbors(self) -> _JsonDict:
|
||
raise NotImplementedError("BGP is not managed via Proxmox API")
|
||
|
||
def get_bgp_neighbors_detail(self, neighbor_address: str = "") -> _JsonDict:
|
||
raise NotImplementedError("BGP is not managed via Proxmox API")
|
||
|
||
def get_route_to(
|
||
self, destination: str = "", protocol: str = "", longer: bool = False
|
||
) -> _JsonDict:
|
||
"""Return routing table entries for the given destination."""
|
||
cmd = f"ip route show {destination} 2>/dev/null || true"
|
||
raw = self._exec_ssh_command(cmd)
|
||
routes: _JsonDict = {}
|
||
for line in raw.splitlines():
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
parts = line.split()
|
||
if not parts:
|
||
continue
|
||
prefix = parts[0]
|
||
next_hop = ""
|
||
out_iface = ""
|
||
proto = "static"
|
||
for i, tok in enumerate(parts):
|
||
if tok == "via" and i + 1 < len(parts):
|
||
next_hop = parts[i + 1]
|
||
if tok == "dev" and i + 1 < len(parts):
|
||
out_iface = parts[i + 1]
|
||
if tok == "proto" and i + 1 < len(parts):
|
||
proto = parts[i + 1]
|
||
|
||
if protocol and protocol.lower() not in proto.lower():
|
||
continue
|
||
|
||
routes.setdefault(prefix, []).append(
|
||
{
|
||
"protocol": proto,
|
||
"current_active": True,
|
||
"last_active": True,
|
||
"age": -1,
|
||
"next_hop": next_hop,
|
||
"outgoing_interface": out_iface,
|
||
"selected_next_hop": True,
|
||
"preference": 1,
|
||
"inactive_reason": "",
|
||
"routing_table": "default",
|
||
"protocol_attributes": {},
|
||
}
|
||
)
|
||
return routes
|
||
|
||
def get_optics(self) -> _JsonDict:
|
||
raise NotImplementedError("Optics not available via Proxmox API")
|
||
|
||
def get_probes_config(self) -> _JsonDict:
|
||
raise NotImplementedError
|
||
|
||
def get_probes_results(self) -> _JsonDict:
|
||
raise NotImplementedError
|
||
|
||
def get_firewall_policies(self) -> _JsonDict:
|
||
raise NotImplementedError("Use the Proxmox firewall API directly")
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# VMs and Containers
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_vm_interfaces(
|
||
self, vmid: int, vm_type: str
|
||
) -> tuple[dict[str, _JsonDict], bool, bool]:
|
||
"""Return network interfaces for a single VM or LXC container.
|
||
|
||
Returns a 3-tuple ``(interfaces, agent_running, agent_enabled)``:
|
||
|
||
* ``interfaces`` – dict keyed by interface name, each with
|
||
NAPALM-compatible fields plus ``ipv4``, ``bridge``, and ``tag``.
|
||
* ``agent_running`` – True if the QEMU Guest Agent responded during
|
||
this call (always False for LXC).
|
||
* ``agent_enabled`` – True if the QEMU Guest Agent is enabled in the
|
||
VM's Proxmox config (always False for LXC).
|
||
|
||
LXC : uses ``/nodes/{node}/lxc/{vmid}/interfaces`` + LXC config
|
||
QEMU : tries QEMU guest agent first, falls back to VM config parsing.
|
||
In both paths the VM config is fetched to derive bridge/tag.
|
||
"""
|
||
_NET_MODELS = {"virtio", "e1000", "e1000e", "vmxnet3", "rtl8139", "ne2k_pci"}
|
||
|
||
def _parse_net_entry(val_str: str) -> tuple[str, str, int | None]:
|
||
"""Parse a Proxmox net config value → (mac_upper, bridge, tag|None)."""
|
||
mac = bridge = ""
|
||
tag: int | None = None
|
||
for part in str(val_str).split(","):
|
||
if "=" not in part:
|
||
continue
|
||
k, v = part.split("=", 1)
|
||
k = k.strip().lower()
|
||
if k in _NET_MODELS:
|
||
mac = v.strip()
|
||
elif k == "bridge":
|
||
bridge = v.strip()
|
||
elif k == "tag":
|
||
try:
|
||
tag = int(v.strip())
|
||
except ValueError:
|
||
pass
|
||
return mac.upper() if mac else "", bridge, tag
|
||
|
||
interfaces: dict[str, _JsonDict] = {}
|
||
|
||
if vm_type == "container":
|
||
# Build iface_name → (bridge, tag) from LXC config
|
||
# LXC net entries look like: net0=name=eth0,bridge=vmbr40,tag=40,...
|
||
lxc_net_map: dict[str, tuple[str, int | None]] = {}
|
||
try:
|
||
config = self._node_api().lxc(vmid).config.get() or {}
|
||
net_re = re.compile(r"^net(\d+)$")
|
||
for key, val in config.items():
|
||
if not net_re.match(key):
|
||
continue
|
||
iface_name = ""
|
||
bridge = ""
|
||
tag: int | None = None
|
||
for part in str(val).split(","):
|
||
if "=" not in part:
|
||
continue
|
||
k, v = part.split("=", 1)
|
||
k = k.strip().lower()
|
||
if k == "name":
|
||
iface_name = v.strip()
|
||
elif k == "bridge":
|
||
bridge = v.strip()
|
||
elif k == "tag":
|
||
try:
|
||
tag = int(v.strip())
|
||
except ValueError:
|
||
pass
|
||
if iface_name:
|
||
lxc_net_map[iface_name] = (bridge, tag)
|
||
except Exception as exc:
|
||
logger.debug("get_vm_interfaces: LXC %s config failed: %s", vmid, exc)
|
||
|
||
try:
|
||
for iface in (self._node_api().lxc(vmid).interfaces.get() or []):
|
||
name = iface.get("name", "")
|
||
if not name or name == "lo":
|
||
continue
|
||
mac = iface.get("hwaddr", "")
|
||
ipv4 = ""
|
||
inet = iface.get("inet", "")
|
||
if inet:
|
||
ipv4 = inet.split("/")[0]
|
||
bridge, tag = lxc_net_map.get(name, ("", None))
|
||
interfaces[name] = {
|
||
"is_up": True,
|
||
"is_enabled": True,
|
||
"description": bridge,
|
||
"mac_address": mac.upper() if mac else "",
|
||
"speed": -1,
|
||
"mtu": 1500,
|
||
"last_flapped": -1.0,
|
||
"ipv4": ipv4,
|
||
"bridge": bridge,
|
||
"tag": tag,
|
||
}
|
||
except Exception as exc:
|
||
logger.debug("get_vm_interfaces: LXC %s ifaces failed: %s", vmid, exc)
|
||
|
||
# LXC containers do not use QEMU Guest Agent
|
||
return interfaces, False, False
|
||
|
||
else:
|
||
# QEMU: pre-fetch VM config to build MAC → (bridge, tag) map and
|
||
# to check whether the QEMU Guest Agent is enabled.
|
||
mac_to_net: dict[str, tuple[str, int | None]] = {} # mac_upper → (bridge, tag)
|
||
net_idx_map: dict[str, tuple[str, str, int | None]] = {} # "netN" → (mac, bridge, tag)
|
||
agent_enabled = False
|
||
try:
|
||
config = self._node_api().qemu(vmid).config.get() or {}
|
||
# Proxmox stores the agent setting as agent=1, agent=0, or
|
||
# agent=enabled=1[,fstrim_cloned_disks=1,...]
|
||
raw_agent = str(config.get("agent", "0"))
|
||
# Treat any truthy value ("1", "enabled=1", ...) as enabled
|
||
agent_enabled = bool(
|
||
raw_agent.strip() in ("1", "true")
|
||
or raw_agent.startswith("enabled=1")
|
||
or raw_agent.startswith("1,")
|
||
)
|
||
net_re = re.compile(r"^net(\d+)$")
|
||
for key, val in config.items():
|
||
m = net_re.match(key)
|
||
if not m:
|
||
continue
|
||
mac, bridge, tag = _parse_net_entry(val)
|
||
iface_key = f"net{m.group(1)}"
|
||
net_idx_map[iface_key] = (mac, bridge, tag)
|
||
if mac:
|
||
mac_to_net[mac] = (bridge, tag)
|
||
except Exception as exc:
|
||
logger.debug("get_vm_interfaces: QEMU %s config fetch failed: %s", vmid, exc)
|
||
|
||
# Try guest agent first
|
||
agent_ok = False
|
||
try:
|
||
agent_result = self._node_api().qemu(vmid).agent("network-get-interfaces").get()
|
||
for iface in (agent_result or {}).get("result", []):
|
||
name = iface.get("name", "")
|
||
if not name or name == "lo":
|
||
continue
|
||
mac = (iface.get("hardware-address", "") or "").upper()
|
||
ipv4 = ""
|
||
for addr in iface.get("ip-addresses", []):
|
||
if addr.get("ip-address-type") == "ipv4":
|
||
ipv4 = addr.get("ip-address", "")
|
||
break
|
||
bridge, tag = mac_to_net.get(mac, ("", None))
|
||
interfaces[name] = {
|
||
"is_up": True,
|
||
"is_enabled": True,
|
||
"description": bridge,
|
||
"mac_address": mac,
|
||
"speed": -1,
|
||
"mtu": 1500,
|
||
"last_flapped": -1.0,
|
||
"ipv4": ipv4,
|
||
"bridge": bridge,
|
||
"tag": tag,
|
||
}
|
||
agent_ok = bool(interfaces)
|
||
except Exception:
|
||
pass
|
||
|
||
if not agent_ok:
|
||
# Fall back to config-only (gives MAC + bridge + tag, no IP)
|
||
for iface_key, (mac, bridge, tag) in net_idx_map.items():
|
||
interfaces[iface_key] = {
|
||
"is_up": False,
|
||
"is_enabled": True,
|
||
"description": bridge,
|
||
"mac_address": mac,
|
||
"speed": -1,
|
||
"mtu": 1500,
|
||
"last_flapped": -1.0,
|
||
"ipv4": "",
|
||
"bridge": bridge,
|
||
"tag": tag,
|
||
}
|
||
|
||
return interfaces, agent_ok, agent_enabled
|
||
|
||
def get_vms(self) -> list[_JsonDict]:
|
||
"""Return all VMs (QEMU) and containers (LXC) on this node.
|
||
|
||
Each entry contains:
|
||
* vmid (int) - Proxmox VM/container ID
|
||
* name (str) - display name
|
||
* type (str) - ``"vm"`` or ``"container"``
|
||
* status (str) - ``"running"``, ``"stopped"``, etc.
|
||
* vcpus (int) - allocated vCPUs
|
||
* memory (int) - configured RAM in megabytes
|
||
* cpu_usage (float) - current CPU utilisation 0.0–1.0 (from last stats cycle)
|
||
* memory_usage (int)- current RSS in megabytes
|
||
* uptime (int) - uptime in seconds (0 if stopped)
|
||
* node (str) - cluster node name
|
||
* interfaces (dict) - network interfaces (NAPALM format + ipv4 field)
|
||
* ipv4 (str) - primary IPv4 address (empty string if unknown)
|
||
"""
|
||
result: list[_JsonDict] = []
|
||
|
||
# QEMU VMs
|
||
try:
|
||
for vm in (self._node_api().qemu.get() or []):
|
||
vmid = int(vm.get("vmid", 0))
|
||
name = vm.get("name", f"vm-{vmid}")
|
||
status = vm.get("status", "unknown")
|
||
cpu_usage = float(vm.get("cpu", 0.0) or 0.0)
|
||
uptime = int(vm.get("uptime", 0) or 0)
|
||
|
||
# mem/maxmem are in bytes
|
||
maxmem_bytes = int(vm.get("maxmem", 0) or 0)
|
||
mem_bytes = int(vm.get("mem", 0) or 0)
|
||
memory_mb = maxmem_bytes // (1024 * 1024)
|
||
memory_usage_mb = mem_bytes // (1024 * 1024)
|
||
|
||
# vcpus can be in "cpus" key for running VMs
|
||
vcpus = int(vm.get("cpus", vm.get("vcpus", 0)) or 0)
|
||
|
||
interfaces, agent_running, agent_enabled = self.get_vm_interfaces(vmid, "vm")
|
||
ipv4 = next(
|
||
(iface["ipv4"] for iface in interfaces.values() if iface.get("ipv4")),
|
||
"",
|
||
)
|
||
|
||
result.append({
|
||
"vmid": vmid,
|
||
"name": name,
|
||
"type": "vm",
|
||
"status": status,
|
||
"vcpus": vcpus,
|
||
"memory": memory_mb,
|
||
"cpu_usage": round(cpu_usage, 4),
|
||
"memory_usage": memory_usage_mb,
|
||
"uptime": uptime,
|
||
"node": self._node_name,
|
||
"interfaces": interfaces,
|
||
"ipv4": ipv4,
|
||
"agent_enabled": agent_enabled,
|
||
"agent_running": agent_running,
|
||
})
|
||
except Exception as exc:
|
||
logger.warning("get_vms: failed to list QEMU VMs: %s", exc)
|
||
|
||
# LXC containers
|
||
try:
|
||
for ct in (self._node_api().lxc.get() or []):
|
||
vmid = int(ct.get("vmid", 0))
|
||
name = ct.get("name", f"ct-{vmid}")
|
||
status = ct.get("status", "unknown")
|
||
cpu_usage = float(ct.get("cpu", 0.0) or 0.0)
|
||
uptime = int(ct.get("uptime", 0) or 0)
|
||
|
||
maxmem_bytes = int(ct.get("maxmem", 0) or 0)
|
||
mem_bytes = int(ct.get("mem", 0) or 0)
|
||
memory_mb = maxmem_bytes // (1024 * 1024)
|
||
memory_usage_mb = mem_bytes // (1024 * 1024)
|
||
|
||
vcpus = int(ct.get("cpus", 0) or 0)
|
||
|
||
interfaces, agent_running, agent_enabled = self.get_vm_interfaces(vmid, "container")
|
||
ipv4 = next(
|
||
(iface["ipv4"] for iface in interfaces.values() if iface.get("ipv4")),
|
||
"",
|
||
)
|
||
|
||
result.append({
|
||
"vmid": vmid,
|
||
"name": name,
|
||
"type": "container",
|
||
"status": status,
|
||
"vcpus": vcpus,
|
||
"memory": memory_mb,
|
||
"cpu_usage": round(cpu_usage, 4),
|
||
"memory_usage": memory_usage_mb,
|
||
"uptime": uptime,
|
||
"node": self._node_name,
|
||
"interfaces": interfaces,
|
||
"ipv4": ipv4,
|
||
"agent_enabled": agent_enabled,
|
||
"agent_running": agent_running,
|
||
})
|
||
except Exception as exc:
|
||
logger.warning("get_vms: failed to list LXC containers: %s", exc)
|
||
|
||
return sorted(result, key=lambda x: x["vmid"])
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# VM power management
|
||
# ------------------------------------------------------------------ #
|
||
|
||
_POWER_ACTIONS_VM = {'start', 'stop', 'shutdown', 'reboot', 'reset'}
|
||
_POWER_ACTIONS_CT = {'start', 'stop', 'shutdown', 'reboot'}
|
||
|
||
def power_vm(self, vmid: int, vm_type: str, action: str) -> dict:
|
||
"""Send a power action to a VM or container on this node.
|
||
|
||
Supported actions for VMs: start, stop, shutdown, reboot, reset
|
||
Supported actions for containers: start, stop, shutdown, reboot
|
||
"""
|
||
allowed = self._POWER_ACTIONS_VM if vm_type == 'vm' else self._POWER_ACTIONS_CT
|
||
if action not in allowed:
|
||
return {"success": False, "error": f"Action '{action}' not supported for {vm_type} (allowed: {sorted(allowed)})"}
|
||
try:
|
||
vm_api = self._node_api().qemu(vmid) if vm_type == 'vm' else self._node_api().lxc(vmid)
|
||
task_id = getattr(vm_api.status, action).post()
|
||
return {"success": True, "task_id": task_id or ""}
|
||
except Exception as exc:
|
||
return {"success": False, "error": str(exc)}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Packages (Debian APT)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_packages(self) -> list[_JsonDict]:
|
||
"""Return installed Debian packages with available-update info.
|
||
|
||
Installed list comes from ``dpkg-query`` via SSH (the Proxmox API
|
||
``/apt/installed`` endpoint is not implemented on PVE 8.x).
|
||
Available updates come from the Proxmox API ``/apt/update``.
|
||
"""
|
||
# Available updates from Proxmox API (keyed by package name)
|
||
upgradable: dict[str, str] = {}
|
||
try:
|
||
for upd in self._api.nodes(self._node_name).apt.update.get():
|
||
pkg = upd.get("Package", "")
|
||
if pkg:
|
||
upgradable[pkg] = upd.get("Version", "")
|
||
except Exception:
|
||
pass
|
||
|
||
# Installed packages via SSH dpkg-query
|
||
raw = self._exec_ssh_command(
|
||
"dpkg-query -W -f='${Package}\\t${Version}\\t${db:Status-Status}\\t${Installed-Size}\\n'"
|
||
" 2>/dev/null"
|
||
)
|
||
result: list[_JsonDict] = []
|
||
for line in raw.splitlines():
|
||
parts = line.strip().split("\t")
|
||
if len(parts) < 2:
|
||
continue
|
||
name = parts[0]
|
||
version = parts[1] if len(parts) > 1 else ""
|
||
status = parts[2] if len(parts) > 2 else "installed"
|
||
size_kb = parts[3] if len(parts) > 3 else "0"
|
||
if not name or status != "installed":
|
||
continue
|
||
size_bytes = int(size_kb) * 1024 if size_kb.isdigit() else 0
|
||
result.append({
|
||
"name": name,
|
||
"version": version,
|
||
"installed": True,
|
||
"description": "",
|
||
"size": size_bytes,
|
||
"source": "pve",
|
||
"upgrade_version": upgradable.get(name, ""),
|
||
})
|
||
return result
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Device warnings
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_device_warnings(self) -> list[_JsonDict]:
|
||
"""Return warnings for the Proxmox node.
|
||
|
||
Currently detects:
|
||
- Available package updates (via Proxmox APT API)
|
||
- Missing / invalid subscription
|
||
"""
|
||
warnings: list[_JsonDict] = []
|
||
|
||
# 1. Available package updates
|
||
try:
|
||
updates = self._api.nodes(self._node_name).apt.update.get()
|
||
if updates:
|
||
warnings.append({
|
||
"code": "updates_available",
|
||
"severity": "info",
|
||
"action": None,
|
||
"meta": {
|
||
"count": len(updates),
|
||
"packages": [u.get("Package", "") for u in updates[:10]],
|
||
},
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
# 2. Subscription status
|
||
try:
|
||
sub = self._get_node_subscription()
|
||
status = sub.get("status", "")
|
||
if status in ("NotFound", "Invalid", "Expired"):
|
||
warnings.append({
|
||
"code": "no_subscription",
|
||
"severity": "warning",
|
||
"action": None,
|
||
"meta": {"status": status},
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
return warnings
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Services (systemd)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_services(self) -> list[_JsonDict]:
|
||
"""Return systemd services with running and enabled state.
|
||
|
||
Uses two ``systemctl`` invocations combined in a single SSH command:
|
||
- ``list-unit-files`` for the static enabled/disabled state
|
||
- ``list-units`` for the live running state
|
||
"""
|
||
raw = self._exec_ssh_command(
|
||
"{ systemctl list-unit-files --type=service --no-pager --no-legend --full 2>/dev/null;"
|
||
" echo '---UNITS---';"
|
||
" systemctl list-units --type=service --all --no-pager --no-legend --full 2>/dev/null;"
|
||
" } || true"
|
||
)
|
||
|
||
# Parse enabled state from list-unit-files
|
||
enabled_map: dict[str, bool] = {}
|
||
section = "files"
|
||
for line in raw.splitlines():
|
||
if line.strip() == "---UNITS---":
|
||
section = "units"
|
||
continue
|
||
parts = line.strip().split(None, 1)
|
||
if len(parts) < 1:
|
||
continue
|
||
unit = parts[0].lstrip("●").strip()
|
||
if not unit.endswith(".service"):
|
||
continue
|
||
name = unit[: -len(".service")]
|
||
if section == "files":
|
||
state = parts[1].strip() if len(parts) > 1 else ""
|
||
enabled_map[name] = state in ("enabled", "enabled-runtime", "static")
|
||
|
||
# Parse running state from list-units
|
||
running_map: dict[str, bool] = {}
|
||
section = "files"
|
||
for line in raw.splitlines():
|
||
if line.strip() == "---UNITS---":
|
||
section = "units"
|
||
continue
|
||
if section != "units":
|
||
continue
|
||
parts = line.strip().lstrip("●").strip().split(None, 4)
|
||
if len(parts) < 4:
|
||
continue
|
||
unit = parts[0]
|
||
if not unit.endswith(".service"):
|
||
continue
|
||
name = unit[: -len(".service")]
|
||
sub_state = parts[3]
|
||
running_map[name] = sub_state == "running"
|
||
|
||
all_names = sorted(set(enabled_map) | set(running_map))
|
||
return [
|
||
{
|
||
"name": name,
|
||
"running": running_map.get(name, False),
|
||
"enabled": enabled_map.get(name, False),
|
||
"pid": 0,
|
||
}
|
||
for name in all_names
|
||
]
|
||
|
||
def manage_service(self, name: str, action: str) -> _JsonDict:
|
||
"""Start / stop / restart / enable / disable a systemd service."""
|
||
if not re.match(r'^[a-zA-Z0-9_\-\.@]+$', name):
|
||
raise ValueError(f"Invalid service name: {name!r}")
|
||
if action not in ('start', 'stop', 'restart', 'enable', 'disable'):
|
||
raise ValueError(f"Invalid action: {action!r}")
|
||
output = self._exec_ssh_command(f"systemctl {action} {name}.service 2>&1 || true")
|
||
return {"success": True, "output": output}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# Available updates
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_available_updates(self) -> list[_JsonDict]:
|
||
"""Return list of upgradable packages from the Proxmox APT API."""
|
||
updates: list[_JsonDict] = []
|
||
try:
|
||
for upd in self._api.nodes(self._node_name).apt.update.get():
|
||
pkg = upd.get("Package", "")
|
||
if not pkg:
|
||
continue
|
||
updates.append({
|
||
"name": pkg,
|
||
"current_version": upd.get("OldVersion", ""),
|
||
"new_version": upd.get("Version", ""),
|
||
})
|
||
except Exception:
|
||
pass
|
||
return sorted(updates, key=lambda u: u["name"])
|
||
|
||
def apply_updates(self, packages: list[str]) -> _JsonDict:
|
||
"""Upgrade the given packages via ``apt-get install`` over SSH."""
|
||
for pkg in packages:
|
||
if not re.match(r'^[a-zA-Z0-9_\-\+\.]+$', pkg):
|
||
raise ValueError(f"Invalid package name: {pkg!r}")
|
||
pkg_args = " ".join(packages)
|
||
output = self._exec_ssh_command(
|
||
f"DEBIAN_FRONTEND=noninteractive apt-get install --only-upgrade -y {pkg_args} 2>&1 || true"
|
||
)
|
||
return {"success": True, "output": output}
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# VLAN provisioning via SDN VNets (OVS-backed nodes only)
|
||
# ------------------------------------------------------------------ #
|
||
|
||
# Proxmox-internal / runtime virtual interface name prefixes that should
|
||
# never be considered physical switch uplinks.
|
||
_VIRTUAL_IFACE_PREFIXES = (
|
||
"fwpr", # Proxmox firewall proxy veth
|
||
"fwln", # Proxmox firewall line veth
|
||
"tap", # VM tap devices
|
||
"veth", # generic veth pairs
|
||
"virbr", # libvirt bridges
|
||
"docker", # Docker virtual interfaces
|
||
"lxcbr", # LXC bridges
|
||
)
|
||
|
||
def _is_physical_uplink(self, iface_name: str, network: dict) -> bool:
|
||
"""Return True if *iface_name* is a physical Ethernet port usable as uplink.
|
||
|
||
Rules:
|
||
- Must not match any known virtual interface name prefix.
|
||
- Must appear in the Proxmox node network config (runtime-only virtual
|
||
interfaces such as ``fwpr*`` or ``tap*`` will not be listed there).
|
||
- Must have a physical-compatible type:
|
||
- ``"eth"`` — regular physical NIC
|
||
- ``"OVSPort"`` — physical NIC attached directly to an OVS bridge
|
||
- ``""`` — untyped (e.g. OVS bond slave, still physical)
|
||
"""
|
||
if any(iface_name.startswith(p) for p in self._VIRTUAL_IFACE_PREFIXES):
|
||
return False
|
||
iface_info = network.get(iface_name)
|
||
if iface_info is None:
|
||
# Not in Proxmox network config → runtime virtual interface
|
||
return False
|
||
return iface_info.get("type", "") in ("eth", "OVSPort", "")
|
||
|
||
def _find_switch_uplink(self) -> str | None:
|
||
"""Return the name of the physical interface connected to a switch.
|
||
|
||
Detection order:
|
||
1. LLDP detailed: physical port whose neighbour advertises Bridge
|
||
capability.
|
||
2. LLDP basic fallback: first physical port with any LLDP neighbour.
|
||
"""
|
||
network = {
|
||
iface["iface"]: iface
|
||
for iface in self._get_node_network()
|
||
if iface.get("iface")
|
||
}
|
||
|
||
# Prefer neighbours that announce Bridge capability
|
||
try:
|
||
for iface_name, neighbour_list in self.get_lldp_neighbors_detail().items():
|
||
if not self._is_physical_uplink(iface_name, network):
|
||
continue
|
||
for nb in neighbour_list:
|
||
caps = nb.get("remote_system_capab", [])
|
||
if any("bridge" in str(c).lower() for c in caps):
|
||
return iface_name
|
||
except Exception:
|
||
pass
|
||
|
||
# Fallback: first physical port with any LLDP neighbour
|
||
try:
|
||
for iface_name, neighbour_list in self.get_lldp_neighbors().items():
|
||
if self._is_physical_uplink(iface_name, network) and neighbour_list:
|
||
return iface_name
|
||
except Exception:
|
||
pass
|
||
|
||
return None
|
||
|
||
def _get_ovs_bridge_for_port(self, port_name: str) -> str | None:
|
||
"""Return the OVS bridge name that *port_name* belongs to, or ``None``.
|
||
|
||
Checks (in order):
|
||
1. Port listed in an OVSBridge's ``ovs_ports``.
|
||
2. Port is a slave of an OVSBond which has an ``ovs_bridge`` reference.
|
||
3. Port itself carries an ``ovs_bridge`` field.
|
||
"""
|
||
network = self._get_node_network()
|
||
by_name: dict[str, _JsonDict] = {
|
||
iface["iface"]: iface for iface in network if iface.get("iface")
|
||
}
|
||
|
||
for iface in network:
|
||
if iface.get("type") == "OVSBridge":
|
||
ports = (iface.get("ovs_ports") or "").split()
|
||
if port_name in ports:
|
||
return iface["iface"]
|
||
|
||
for iface in network:
|
||
if iface.get("type") == "OVSBond":
|
||
slaves = (iface.get("slaves") or "").split()
|
||
if port_name in slaves:
|
||
bridge = iface.get("ovs_bridge", "")
|
||
if bridge:
|
||
return bridge
|
||
|
||
port_info = by_name.get(port_name, {})
|
||
return port_info.get("ovs_bridge") or None
|
||
|
||
def _is_cluster_master(self) -> bool:
|
||
"""Return ``True`` if this node is the Corosync quorum coordinator.
|
||
|
||
The coordinator is the online cluster node with the lowest ``nodeid``.
|
||
On standalone (non-clustered) nodes this always returns ``True``.
|
||
"""
|
||
try:
|
||
status = self._api.cluster.status.get() or []
|
||
node_entries = [e for e in status if e.get("type") == "node"]
|
||
if not node_entries:
|
||
return True # Standalone node — no cluster
|
||
online_nodes = [n for n in node_entries if n.get("online", 0)]
|
||
if not online_nodes:
|
||
return True # All nodes offline → assume we're the master
|
||
min_id = min(int(n.get("nodeid", 9999)) for n in online_nodes)
|
||
for n in online_nodes:
|
||
if (
|
||
n.get("name") == self._node_name
|
||
and int(n.get("nodeid", 9999)) == min_id
|
||
):
|
||
return True
|
||
return False
|
||
except Exception:
|
||
return True # Can't determine → assume standalone, proceed
|
||
|
||
def _get_sdn_zone_for_bridge(self, bridge_name: str) -> str | None:
|
||
"""Return the SDN zone ID whose ``bridge`` field matches *bridge_name*.
|
||
|
||
In Proxmox SDN each zone is linked to exactly one OVS bridge via the
|
||
``bridge`` property. We look for that mapping so the VNet is always
|
||
created in the correct zone instead of guessing by type order.
|
||
|
||
Falls back to the first zone if no bridge match is found.
|
||
"""
|
||
try:
|
||
zones = self._get_sdn_zones()
|
||
# Exact match on bridge field
|
||
for zone in zones:
|
||
if zone.get("bridge") == bridge_name:
|
||
return zone.get("zone")
|
||
# Fallback: first available zone
|
||
if zones:
|
||
return zones[0].get("zone")
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
def set_vlan(self, vlan_id: int, config) -> None:
|
||
"""Create (or update) a VLAN via an SDN VNet on this Proxmox node.
|
||
|
||
Pre-flight checks (all must pass to proceed):
|
||
|
||
1. Finds the physical uplink port connected to a switch via LLDP.
|
||
Physical ports are those with type ``eth``, ``OVSPort``, or ``""``
|
||
in the Proxmox network config (excludes runtime virtuals like
|
||
``fwpr*``, ``tap*``, etc.).
|
||
2. Verifies that uplink is part of an OVS bridge or OVS bond.
|
||
3. Confirms this node is the Corosync quorum master (lowest node-id).
|
||
Non-master nodes return silently — the master handles VNet creation.
|
||
|
||
The SDN VNet is named ``vlan{vid:04d}`` (e.g. ``vlan0007`` for VID 7).
|
||
If the VNet already exists its alias is updated. After creating /
|
||
updating the VNet the SDN configuration is reloaded via
|
||
``PUT /cluster/sdn``.
|
||
|
||
Args:
|
||
vlan_id: VLAN identifier (1-4094).
|
||
config: Dict that may contain ``"name"`` for the VLAN alias.
|
||
"""
|
||
name: str = (
|
||
(config.get("name") or f"VLAN{vlan_id}") if config else f"VLAN{vlan_id}"
|
||
)
|
||
vnet_id = f"vlan{vlan_id:04d}"
|
||
|
||
# 1. Find switch uplink via LLDP
|
||
uplink = self._find_switch_uplink()
|
||
if uplink is None:
|
||
raise ConnectionException(
|
||
f"set_vlan({vlan_id}): no LLDP-detected switch uplink found"
|
||
f" on node {self._node_name!r}"
|
||
)
|
||
|
||
# 2. Verify uplink is part of an OVS bridge
|
||
ovs_bridge = self._get_ovs_bridge_for_port(uplink)
|
||
if ovs_bridge is None:
|
||
raise ConnectionException(
|
||
f"set_vlan({vlan_id}): uplink {uplink!r} is not part of an OVS bridge"
|
||
)
|
||
|
||
# 3. Only the Corosync master manages SDN VNets
|
||
if not self._is_cluster_master():
|
||
return
|
||
|
||
# 4. Find the SDN zone linked to this OVS bridge
|
||
zone = self._get_sdn_zone_for_bridge(ovs_bridge)
|
||
if not zone:
|
||
raise ConnectionException(
|
||
f"set_vlan({vlan_id}): no SDN zone found for bridge {ovs_bridge!r}"
|
||
)
|
||
|
||
# 5. Create or update the VNet
|
||
try:
|
||
self._api.cluster.sdn.vnets.post(
|
||
vnet=vnet_id,
|
||
zone=zone,
|
||
tag=vlan_id,
|
||
alias=name,
|
||
)
|
||
except ResourceException as exc:
|
||
err_str = str(exc).lower()
|
||
if "already exists" in err_str or "duplicate" in err_str or "500" in err_str:
|
||
# VNet already exists — update alias
|
||
try:
|
||
self._api.cluster.sdn.vnets(vnet_id).put(alias=name)
|
||
except Exception:
|
||
pass # Best effort
|
||
else:
|
||
raise ConnectionException(
|
||
f"set_vlan({vlan_id}): failed to create VNet {vnet_id!r}: {exc}"
|
||
) from exc
|
||
|
||
# 6. Reload SDN config so the VNet becomes active
|
||
try:
|
||
self._api.cluster.sdn.put()
|
||
except Exception:
|
||
pass # Best effort — may not be needed on older PVE versions
|
||
|
||
def delete_vlan(self, vlan_id: int) -> None:
|
||
"""Delete the SDN VNet corresponding to *vlan_id*.
|
||
|
||
The VNet is identified by the canonical name ``vlan{vid:04d}``.
|
||
Only the Corosync quorum master performs the deletion — non-master
|
||
nodes return silently.
|
||
|
||
After deletion the SDN configuration is reloaded via
|
||
``PUT /cluster/sdn``.
|
||
|
||
Args:
|
||
vlan_id: VLAN identifier to delete.
|
||
"""
|
||
if not self._is_cluster_master():
|
||
return # Non-master node — master handles VNet deletion
|
||
|
||
vnet_id = f"vlan{vlan_id:04d}"
|
||
try:
|
||
self._api.cluster.sdn.vnets(vnet_id).delete()
|
||
except ResourceException as exc:
|
||
err_str = str(exc).lower()
|
||
if "does not exist" in err_str or "404" in str(exc):
|
||
return # Already gone — treat as success
|
||
raise ConnectionException(
|
||
f"delete_vlan({vlan_id}): failed to delete VNet {vnet_id!r}: {exc}"
|
||
) from exc
|
||
|
||
# Reload SDN config
|
||
try:
|
||
self._api.cluster.sdn.put()
|
||
except Exception:
|
||
pass # Best effort
|
||
|
||
# ------------------------------------------------------------------ #
|
||
# SNMP / Health
|
||
# ------------------------------------------------------------------ #
|
||
|
||
def get_snmp_config(self):
|
||
"""Return SNMP agent config if snmpd is installed and running on the node.
|
||
|
||
Uses _exec_ssh_command (Proxmox API exec or SSH) to inspect the node.
|
||
Returns a SNMPConfigDict or None.
|
||
"""
|
||
try:
|
||
from napalm_device_types.models import SNMPConfigDict
|
||
except ImportError:
|
||
return None
|
||
|
||
running = (
|
||
self._exec_ssh_command("systemctl is-active snmpd 2>/dev/null || true").strip()
|
||
== "active"
|
||
)
|
||
if not running:
|
||
return None
|
||
|
||
community = "public"
|
||
port = 161
|
||
try:
|
||
conf = self._exec_ssh_command(
|
||
"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") and len(parts) >= 2:
|
||
community = parts[1]
|
||
break
|
||
elif kw == "com2sec" and len(parts) >= 4:
|
||
community = parts[3]
|
||
break
|
||
except Exception:
|
||
pass
|
||
|
||
return SNMPConfigDict(running=True, community=community, port=port, version="2c")
|
||
|
||
def run_device_action(self, action: str) -> dict:
|
||
"""Execute a named administrative action on the Proxmox node."""
|
||
if action == "fix_snmp":
|
||
return self._action_fix_snmp()
|
||
raise NotImplementedError(f"Unknown action: {action!r}")
|
||
|
||
def _action_fix_snmp(self) -> dict:
|
||
"""Install, configure and start snmpd on the Proxmox node.
|
||
|
||
Proxmox runs Debian/Linux underneath. _exec_ssh_command runs as root
|
||
(either via Proxmox API execute endpoint or SSH with root credentials),
|
||
so no sudo is needed.
|
||
"""
|
||
import base64 as _b64
|
||
lines: [].__class__ = []
|
||
|
||
# 1. Install snmpd and snmp client tools
|
||
install_out = self._exec_ssh_command(
|
||
"DEBIAN_FRONTEND=noninteractive apt-get install -y snmpd snmp 2>&1 | tail -5"
|
||
)
|
||
lines.append(f"[install] {install_out.strip()[-200:]}")
|
||
|
||
# 2. Detect the IP this connection comes from (for firewall rule)
|
||
netork_ip = ""
|
||
try:
|
||
raw = self._exec_ssh_command(
|
||
"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
|
||
|
||
# 3. Write snmpd.conf via /tmp (no permission issues)
|
||
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._exec_ssh_command(f"echo {conf_b64} | base64 -d > /tmp/netork_snmpd.conf")
|
||
self._exec_ssh_command(
|
||
"mv /tmp/netork_snmpd.conf /etc/snmp/snmpd.conf && "
|
||
"chown root:root /etc/snmp/snmpd.conf && chmod 644 /etc/snmp/snmpd.conf"
|
||
)
|
||
verify = self._exec_ssh_command("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]}")
|
||
|
||
# 4. Open firewall if ufw is present
|
||
if netork_ip:
|
||
try:
|
||
ufw = self._exec_ssh_command("command -v ufw 2>/dev/null").strip()
|
||
if ufw:
|
||
parts = netork_ip.rsplit(".", 1)
|
||
subnet = f"{parts[0]}.0/24" if len(parts) == 2 else netork_ip
|
||
fw_out = self._exec_ssh_command(
|
||
f"ufw allow from {subnet} to any port 161 proto udp 2>&1"
|
||
)
|
||
lines.append(f"[firewall/ufw] {fw_out.strip()[:200]}")
|
||
except Exception as exc:
|
||
lines.append(f"[firewall] skipped — {exc}")
|
||
|
||
# 5. Stop and restart snmpd cleanly (no DBus needed for stop+start)
|
||
self._exec_ssh_command(
|
||
"service snmpd stop 2>/dev/null; pkill -9 snmpd 2>/dev/null; true"
|
||
)
|
||
import time as _time
|
||
_time.sleep(1)
|
||
start_out = self._exec_ssh_command(
|
||
"service snmpd start 2>&1 || systemctl start snmpd 2>&1 || true"
|
||
)
|
||
lines.append(f"[service] {start_out.strip()[-200:]}")
|
||
|
||
# 6. Verify via local probe
|
||
_time.sleep(2)
|
||
probe_out = self._exec_ssh_command(
|
||
"snmpget -v2c -cpublic -t2 -r0 -Ov 127.0.0.1 1.3.6.1.2.1.1.1.0 2>&1 || true"
|
||
).strip()
|
||
_snmp_types = ("STRING:", "INTEGER:", "OID:", "Timeticks:", "Hex-STRING:", "IpAddress:")
|
||
success = any(t in probe_out for t in _snmp_types)
|
||
if success:
|
||
lines.append("[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 get_disk_smart(self) -> dict:
|
||
"""Return SMART health data for all disks on this Proxmox node.
|
||
|
||
Combines the disk list (model, size, wearout) with per-disk SMART
|
||
data (health, temperature, percentage used).
|
||
|
||
Returns a dict keyed by device path, e.g. {"/dev/nvme0n1": {...}}.
|
||
"""
|
||
import re as _re
|
||
|
||
result: dict = {}
|
||
try:
|
||
disks = self._node_api().disks.list.get() or []
|
||
except Exception:
|
||
return result
|
||
|
||
for disk in disks:
|
||
dev = disk.get("devpath") or disk.get("dev")
|
||
if not dev:
|
||
continue
|
||
entry: dict = {
|
||
"model": disk.get("model", ""),
|
||
"serial": disk.get("serial", ""),
|
||
"type": disk.get("type", ""),
|
||
"size": disk.get("size", 0),
|
||
"health": disk.get("health", "unknown").lower(),
|
||
"wearout": disk.get("wearout"), # NVMe wear indicator 0-100
|
||
"temperature": None,
|
||
"percentage_used": None,
|
||
"available_spare": None,
|
||
"reallocated_sectors": None,
|
||
"power_on_hours": None,
|
||
}
|
||
try:
|
||
smart = self._node_api().disks.smart.get(disk=dev) or {}
|
||
# health from SMART endpoint may be more accurate
|
||
if smart.get("health"):
|
||
entry["health"] = smart["health"].lower()
|
||
|
||
text = smart.get("text", "")
|
||
# Parse temperature
|
||
m = _re.search(r"Temperature[^:]*:\s*(\d+)\s*Celsius", text)
|
||
if m:
|
||
entry["temperature"] = int(m.group(1))
|
||
# NVMe-specific
|
||
m = _re.search(r"Percentage Used:\s*(\d+)%", text)
|
||
if m:
|
||
entry["percentage_used"] = int(m.group(1))
|
||
m = _re.search(r"Available Spare:\s*(\d+)%", text)
|
||
if m:
|
||
entry["available_spare"] = int(m.group(1))
|
||
m = _re.search(r"Power On Hours:\s*([\d,]+)", text)
|
||
if m:
|
||
entry["power_on_hours"] = int(m.group(1).replace(",", ""))
|
||
# HDD-specific SMART attributes
|
||
for attr in smart.get("attributes", []):
|
||
name = attr.get("name", "").lower()
|
||
raw = attr.get("raw", "")
|
||
try:
|
||
raw_int = int(str(raw).split()[0])
|
||
except (ValueError, TypeError):
|
||
raw_int = None
|
||
if "temperature" in name and raw_int is not None:
|
||
entry["temperature"] = raw_int
|
||
elif "reallocated" in name and "sector" in name and raw_int is not None:
|
||
entry["reallocated_sectors"] = raw_int
|
||
elif "power_on" in name and raw_int is not None:
|
||
entry["power_on_hours"] = raw_int
|
||
except Exception:
|
||
pass
|
||
result[dev] = entry
|
||
|
||
return result
|