Author SHA1 Message Date
christianmanivong 885c7e1f53 feat!: declare the guest agent per guest OS, and share what provisioning needs
CI / test (3.10) (push) Successful in 39s
CI / test (3.11) (push) Successful in 36s
CI / test (3.12) (push) Successful in 39s
CI / test (3.10) (pull_request) Successful in 32s
CI / test (3.11) (pull_request) Successful in 23s
CI / test (3.12) (pull_request) Successful in 24s
GUEST_AGENT_PACKAGES / GUEST_AGENT_RUNCMD described one agent per
hypervisor. A hypervisor that provisions FreeBSD and OpenBSD guests needs
a different agent per guest, and netOrk has to know which guests a driver
can provision at all (NetOrk/netork#794). Both now come from one
class-level mapping netOrk reads off the driver class:

    GUEST_AGENTS = {guest_os: (packages, runcmd)}

Its keys are the guests the driver provisions, and
create_vm_from_cloud_init(guest_os=...) refuses any other. The base keeps
Linux with qemu-guest-agent.

New module `provisioning`, generic for every hypervisor:
- QEMU_GUEST_AGENTS: the QEMU agent for Linux, FreeBSD and OpenBSD, with
  package names and service commands verified on FreeBSD 15.1 and
  OpenBSD 7.9 cloud images (NetOrk/netork#793).
- network_config(): cloud-init's network-config v2 with MAC matching,
  moved here from napalm-vmware. FreeBSD's nuageinit reads no other
  version.
- split_compression(): how a packed image (.xz, .gz, .bz2, .zst) is
  unpacked.

BREAKING CHANGE: GUEST_AGENT_PACKAGES and GUEST_AGENT_RUNCMD are gone;
drivers declare GUEST_AGENTS instead.
2026-10-08 07:53:21 +02:00
christianmanivong eb80d5cb0d Merge pull request 'feat(container-engine): run the engine CLI privileged on request' (#19) from feat/privileged-cli into main
CI / test (3.10) (push) Successful in 23s
CI / test (3.11) (push) Successful in 22s
CI / test (3.12) (push) Successful in 23s
2026-10-07 20:11:13 +00:00
christianmanivong e50e497939 feat(container-engine): run the engine CLI privileged on request
CI / test (3.10) (push) Successful in 23s
CI / test (3.11) (push) Successful in 23s
CI / test (3.12) (push) Successful in 24s
CI / test (3.10) (pull_request) Successful in 23s
CI / test (3.11) (pull_request) Successful in 22s
CI / test (3.12) (pull_request) Successful in 23s
run_cli() and stream_cli() take privileged=, passed to the driver's
run_command()/open_stream(): an engine that refuses the login user can
be retried as root, the way the driver gains root for any command.
netOrk uses it to retry a refused container start/stop/restart with
sudo (NetOrk/netork#773).
2026-10-07 22:08:37 +02:00
christianmanivong 3aed0b48d7 Merge pull request 'feat: a public command channel and container engine access' (#18) from feat/container-engine-channel into main
CI / test (3.10) (push) Successful in 22s
CI / test (3.11) (push) Successful in 21s
CI / test (3.12) (push) Successful in 23s
2026-10-07 16:01:51 +00:00
christianmanivong 735b683028 feat: a public command channel and container engine access
CI / test (3.10) (push) Successful in 25s
CI / test (3.11) (push) Successful in 22s
CI / test (3.12) (push) Successful in 23s
CI / test (3.10) (pull_request) Successful in 24s
CI / test (3.11) (pull_request) Successful in 22s
CI / test (3.12) (pull_request) Successful in 22s
netOrk reached a host's shell through napalm-linux's private `_send`: an
interactive PTY with stdout and stderr merged and no exit code. For Docker
it went further and opened its own paramiko connections around the driver.

This adds the access layer only, no container model (NetOrk/netork#765):

- `channel.py`: `CommandChannelMixin` declares `run_command()` (stdout,
  stderr, exit code) and `open_stream()` (a `ByteStream` to a running
  command). Declared under TYPE_CHECKING, so hasattr stays truthful.
  `run_on_transport()` and `open_stream_on_transport()` implement both on
  a paramiko exec channel for any SSH driver; `ParamikoExecStream` drains
  stderr on every read so the shared window never stalls.
- `container_engine.py`: `ContainerEngineMixin` with `container_engines()`
  and `open_container_engine()`. The returned `ContainerEngineConnection`
  opens the engine API over `<binary> system dial-stdio` and runs the CLI
  with caller-chosen arguments. The driver decides the binary through
  `_container_engine_binary()`; callers never see the path.
- README: why the container model lives in netOrk and not here.

Additive; the Docker*Dict declarations stay until netOrk no longer reads
them. Version 2.6.0.

Refs #17
2026-10-07 17:58:57 +02:00
christianmanivong 08ec32e93a Merge pull request 'ci: run the tests and build the package on every push and pull request' (#16) from ci/workflow into main
CI / test (3.10) (push) Successful in 23s
CI / test (3.11) (push) Successful in 22s
CI / test (3.12) (push) Successful in 24s
2026-10-07 06:22:15 +00:00
christianmanivong bd43bd75fa ci: run the tests and build the package on every push and pull request
CI / test (3.10) (push) Successful in 29s
CI / test (3.11) (push) Successful in 28s
CI / test (3.12) (push) Successful in 27s
CI / test (3.10) (pull_request) Successful in 24s
CI / test (3.11) (pull_request) Successful in 25s
CI / test (3.12) (pull_request) Successful in 26s
The same job as napalm-fritzbox and napalm-opnsense: Python 3.10, 3.11 and
3.12, pytest, then wheel and sdist. napalm-device-types comes from
git.netork.io first, because PyPI has an unrelated package of that name.

This is the package the other drivers install from git first; here it
installs itself.
2026-10-07 07:52:44 +02:00
christianmanivong 1112191aec Merge pull request 'feat: read the sockets of a host without ss from netstat, and procd's services' (#15) from feat/listening-netstat into main 2026-10-07 05:21:31 +00:00
christianmanivong b8b89acee1 feat: read the sockets of a host without ss from netstat, and procd's services
OpenWrt has no ss: the listening-socket command fell through to [no-ss]
and the reading raised ListeningSocketsUnavailable. Without ss it now runs
netstat -lntup (busybox and net-tools both), and the cgroups are read for
its PIDs alike.

- netstat names a socket's process as PID/Program; a UDP line has no
  state column, "-" is a socket without a process, and net-tools prints
  tcp6/udp6 and program names with a space ("sshd: /usr/sbin").
- On OpenWrt the cgroup is /services/<name>/<instance>, so the unit is
  the procd service -- a jailed one too, whose PID is not the one procd
  reports (dnsmasq under ujail).
- A host with neither tool still raises ListeningSocketsUnavailable.

Checked against a real OpenWrt 25.12.2 access point, and that a Linux and
a Proxmox host still read the same sockets with the longer command.

2.5.0. For netOrk#673.
2026-10-07 07:20:27 +02:00
christianmanivong f833e23422 Merge pull request 'feat: read what listens on which address, and which service it is, once for every driver' (#7) from feat/listening-sockets into main 2026-10-06 16:19:22 +00:00
christianmanivong b97ec654a0 feat: read what listens on which address, and which service it is, once for every driver
Whether a service is reachable from outside its host is decided by the
address it listens on: 0.0.0.0:5432 is, 127.0.0.1:5432 is not. Reading
that is the same on every Linux host, so the command and its parse live
here once and a driver only carries the command across
(ListeningSocketsMixin, _run_listening_sockets_command).

One framed round trip: ss -lntup for every listening TCP and bound UDP
socket, then /proc/<pid>/cgroup for each process holding one, which names
the systemd service (v2, nested slices, v1's name=systemd hierarchy) or
the container (docker-<id>.scope, /docker/<id>) it runs in.

- Root: only root sees every process. The script goes as one sh -c
  argument, so a sudo -n prefix covers all of it; when that brings no
  report back the reading runs again unprivileged and says it is not
  attributed.
- No -H: iproute2 before 4.10 fails on it, which would read as nothing
  listening. The header is skipped instead.
- A host without ss raises ListeningSocketsUnavailable; a report cut short
  or a failing ss raises ValueError.
- The reading is raw: docker-proxy shows up as docker.service, loopback as
  loopback. What counts as reachable is the consumer's call.

2.4.0. For netOrk#658.
2026-10-06 18:18:56 +02:00
christianmanivong c2d8d4a0d2 Merge pull request 'feat: say where a pending update comes from, whether it is a security fix, and whether the host needs a reboot' (#6) from feat/update-origin-host-status into main 2026-10-05 22:19:25 +00:00
16 changed files with 1675 additions and 24 deletions
+45
View File
@@ -0,0 +1,45 @@
name: CI
on:
push:
branches: ["**"]
pull_request:
branches: ["**"]
jobs:
test:
runs-on: ubuntu-latest
strategy:
fail-fast: false
matrix:
python-version: ["3.10", "3.11", "3.12"]
steps:
- name: Checkout
uses: actions/checkout@v4
- name: Setup Python
uses: actions/setup-python@v5
with:
python-version: ${{ matrix.python-version }}
cache: pip
- name: Install package with dev extras
run: |
python -m pip install --upgrade pip
python -m pip install -e ".[dev]"
- name: Run unit tests
run: |
python -m pytest -q --tb=short
- name: Build wheel and sdist
run: |
python -m pip install build
python -m build
- name: Upload dist artifacts
# v4 refuses to run on Gitea ("not currently supported on GHES").
uses: actions/upload-artifact@v3
with:
name: dist-${{ matrix.python-version }}
path: dist/*
+41
View File
@@ -86,6 +86,16 @@ dnf-automatic). `package_updates` holds the shared apt and dnf parsers: apt's su
an update's `origin`, a `-security` suite makes it a security update, and dnf's security
advisories do the same.
`ListeningSocketsMixin` (`get_listening_sockets`) is mixed in the same way: every listening
TCP and bound UDP socket from `ss -lntup`, with the systemd service or container behind it
from `/proc/<pid>/cgroup`, in one round trip. A driver supplies
`_run_listening_sockets_command(command, privileged=)`; the command arrives as one `sh -c`
argument, so a `sudo -n` prefix covers all of it. Without root `ss` names only the login
user's processes, and the reading says so (`attributed: false`) instead of failing. A host
without `ss` is read with `netstat -lntup` (OpenWrt's busybox, old net-tools); on OpenWrt the
cgroup names the procd service (`/services/<name>/<instance>`). A host with neither raises
`ListeningSocketsUnavailable`.
**Update readers raise when they cannot read.** `get_available_updates` returns an empty
list only when nothing is pending; netOrk keeps "pending since" per package, and an empty
list for "don't know" would reset it.
@@ -99,6 +109,27 @@ trip (`list-unit-files` plus one `systemctl show` over every loaded unit) instea
`is-enabled` and a `show` per unit. A host without systemd raises `SystemdUnavailable`,
a `NotImplementedError`, so a driver can fall back to another init system.
### Access channels: why there is no container model here
`CommandChannelMixin` declares a driver's public channel to its host:
- `run_command(command, *, privileged=False, timeout=60, stdin=None)` returns a `CommandResult(stdout, stderr, exit_code)`;
- `open_stream(command, *, privileged=False)` returns a `ByteStream` to the command's stdin and stdout.
The mixin sits in `DeviceTypeDriver` and only declares, so `hasattr(driver, "run_command")` is true exactly where a driver implements it. For SSH drivers, `channel.run_on_transport` and `channel.open_stream_on_transport` are the implementation over a paramiko exec channel: no PTY, stderr kept apart, a real exit code. How a command gains root stays with the driver.
**`ContainerEngineMixin` is deliberately different from `SystemdServicesMixin`.** The systemd mixin owns its command and its parse; this one owns neither.
- `container_engines()` says which engines the host offers (`[{"engine": "docker", "api": "docker-engine"}]`).
- `open_container_engine(engine)` returns a `ContainerEngineConnection`:
- `open_api()` is a stream to the engine's API (`docker system dial-stdio`);
- `run_cli(args)` and `stream_cli(args)` run the engine's CLI with arguments the caller chooses. With `privileged=True` the driver runs it the way `run_command` gains root, for an engine that refuses the login user.
- The driver decides only *how* the engine is reached. Its hook `_container_engine_binary(engine)` returns, for QNAP, the Container Station path, and the caller never sees it.
**What runs on the engine, and what to do with it, is netOrk's** (NetOrk/netork#765): the container model, the Engine API requests, compose, updates. Above the connection everything is specific to the service, so this package abstracts the connection and nothing more.
A driver mixes `ContainerEngineMixin` in itself; `OSDriver` does not carry it. A Windows host is an OS driver too, and has no `dial-stdio` to offer over WinRM.
**Never half-close early.** `ByteStream.write` does not close anything; only `close_write` does. `dial-stdio` hands the daemon a half-close as "client gone", and the daemon then answers an unfinished request with HTTP 499.
A function class may use the **template form** — public method concrete, the
device-specific part a `_hook` declared under `if TYPE_CHECKING` — *when the base
genuinely does work* on the result: normalising, sorting, validating, or orchestrating
@@ -281,6 +312,16 @@ class ProxmoxDriver(HypervisorDriver):
...
```
A hypervisor that provisions VMs from cloud images declares, per guest OS,
the agent cloud-init installs so the hypervisor can read the new VM's IP:
`GUEST_AGENTS = {guest_os: (packages, runcmd)}`. Its keys are the guests the
driver can provision, and `create_vm_from_cloud_init(guest_os=...)` refuses
any other. A QEMU-based driver takes `provisioning.QEMU_GUEST_AGENTS`
(Linux, FreeBSD, OpenBSD). `provisioning.network_config()` writes the
network-config v2 every cloud-init flavour reads (FreeBSD's nuageinit reads
no other), and `provisioning.split_compression()` says how a packed image is
unpacked; both are the same for every hypervisor.
### OS / Linux
```python
+35
View File
@@ -34,7 +34,9 @@ Role bases -- what a device *is*:
Function classes -- what a device *can do*. Shared behaviour lives here once
instead of being restated on every role that happens to need it:
* :class:`~napalm_device_types.channel.CommandChannelMixin`
* :class:`~napalm_device_types.config_lifecycle.ConfigLifecycleMixin`
* :class:`~napalm_device_types.container_engine.ContainerEngineMixin`
* :class:`~napalm_device_types.dhcp.DhcpServerMixin`
* :class:`~napalm_device_types.firewall_rules.FirewallRuleMixin`
* :class:`~napalm_device_types.health_metrics.HealthMetricsMixin`
@@ -42,6 +44,7 @@ instead of being restated on every role that happens to need it:
* :class:`~napalm_device_types.host_reboot.HostRebootMixin`
* :class:`~napalm_device_types.interface_filter.InterfaceFilterMixin`
* :class:`~napalm_device_types.kernel.KernelFactsMixin`
* :class:`~napalm_device_types.listening.ListeningSocketsMixin`
* :class:`~napalm_device_types.mac_acl.MacAclMixin`
* :class:`~napalm_device_types.nat_vpn.NatVpnMixin`
* :class:`~napalm_device_types.packages.PackageManagementMixin`
@@ -57,7 +60,20 @@ Introspection -- :func:`~napalm_device_types.roles.roles_of`,
from napalm_device_types.base import DeviceTypeDriver, FingerprintRule, PortSpec
from napalm_device_types.access_point import AccessPointDriver
from napalm_device_types.channel import (
ByteStream,
CommandChannelMixin,
CommandResult,
ParamikoExecStream,
open_stream_on_transport,
run_on_transport,
)
from napalm_device_types.config_lifecycle import ConfigLifecycleMixin
from napalm_device_types.container_engine import (
ContainerEngineConnection,
ContainerEngineMixin,
ContainerEngineUnavailable,
)
from napalm_device_types.dhcp import DhcpServerMixin, normalize_cidr, normalize_mac
from napalm_device_types.firewall import FirewallDriver
from napalm_device_types.hypervisor import HypervisorDriver
@@ -68,6 +84,12 @@ from napalm_device_types.host_reboot import HostRebootMixin
from napalm_device_types.interface_filter import InterfaceFilterMixin
from napalm_device_types.kernel import KERNEL_FACTS_COMMAND, KernelFactsMixin, parse_kernel_facts
from napalm_device_types.lag import add_lag_interfaces
from napalm_device_types.listening import (
LISTENING_SOCKETS_COMMAND,
ListeningSocketsMixin,
ListeningSocketsUnavailable,
parse_listening_sockets,
)
from napalm_device_types.mac_acl import MacAclMixin
from napalm_device_types.media import MediaDriver
from napalm_device_types.nat_vpn import NatVpnMixin
@@ -97,6 +119,15 @@ from napalm_device_types.switch import SwitchDriver
__all__ = [
"AccessPointDriver",
"ByteStream",
"CommandChannelMixin",
"CommandResult",
"ContainerEngineConnection",
"ContainerEngineMixin",
"ContainerEngineUnavailable",
"ParamikoExecStream",
"open_stream_on_transport",
"run_on_transport",
"ConfigLifecycleMixin",
"DeviceTypeDriver",
"DhcpServerMixin",
@@ -123,6 +154,10 @@ __all__ = [
"KernelFactsMixin",
"KERNEL_FACTS_COMMAND",
"parse_kernel_facts",
"LISTENING_SOCKETS_COMMAND",
"ListeningSocketsMixin",
"ListeningSocketsUnavailable",
"parse_listening_sockets",
"PhoneDriver",
"PingSweepMixin",
"PortSpec",
+2 -1
View File
@@ -12,6 +12,7 @@ from typing import NamedTuple
from napalm.base import NetworkDriver
from napalm_device_types.channel import CommandChannelMixin
from napalm_device_types.host_reboot import HostRebootMixin
from napalm_device_types.ping_sweep import PingSweepMixin
@@ -48,7 +49,7 @@ class PortSpec(NamedTuple):
mandatory: bool = False
class DeviceTypeDriver(PingSweepMixin, HostRebootMixin, NetworkDriver):
class DeviceTypeDriver(PingSweepMixin, HostRebootMixin, CommandChannelMixin, NetworkDriver):
"""Common base for all netOrk device-type drivers.
Sits between napalm.base.NetworkDriver and the type-specific abstract
+188
View File
@@ -0,0 +1,188 @@
# -*- coding: utf-8 -*-
"""The command channel: a public way to run a command on a host, or open a stream.
A driver already reaches its host, over SSH, a REST API or WinRM. What it did
not offer was a *public* way for the layer above to use that reach: netOrk sent
shell strings through napalm-linux's private ``_send``, an interactive PTY that
merges stdout and stderr and has no exit code.
:class:`CommandChannelMixin` declares two methods, and a driver that can
implements them:
* :meth:`run_command` runs a command and returns stdout, stderr and the exit code.
* :meth:`open_stream` starts a command and returns a :class:`ByteStream` to its
stdin and stdout. That is how netOrk speaks the Docker Engine API, over
``docker system dial-stdio`` (see :mod:`napalm_device_types.container_engine`).
Like the role bases, the mixin only declares (under ``TYPE_CHECKING``), so
``hasattr(driver, "run_command")`` stays a truthful answer.
:func:`run_on_transport` and :func:`open_stream_on_transport` are the SSH
implementation over a paramiko ``Transport``, for any SSH driver: an exec
channel, so no PTY, separate stderr and a real exit status. How a command gains
root (``sudo -S`` with the password on stdin, ``sudo -n``, or nothing as root)
stays with the driver.
"""
from __future__ import annotations
import socket
import time
from typing import Any, NamedTuple, Optional, Protocol, TYPE_CHECKING
#: How much of a stream's stderr is kept. Enough to tell "permission denied"
#: from "no such command"; a stream that writes megabytes there keeps its tail.
STDERR_TAIL = 64 * 1024
_READ = 32 * 1024
_POLL = 0.01
class CommandResult(NamedTuple):
"""What a command printed, and how it ended."""
stdout: str
stderr: str
exit_code: int
class ByteStream(Protocol):
"""A running command's stdin and stdout, as bytes.
``read`` returns ``b""`` at end of stream and raises :class:`TimeoutError`
when nothing arrives in time. ``write`` never closes anything:
``dial-stdio`` hands the daemon a half-close as "client gone", and the
daemon then answers an unfinished request with HTTP 499 (netork#771). Only
``close_write`` closes the writing side.
"""
def read(self, max_bytes: int, timeout: Optional[float] = None) -> bytes: ...
def write(self, data: bytes) -> None: ...
def close_write(self) -> None: ...
def close(self) -> None: ...
@property
def exit_status(self) -> Optional[int]: ...
@property
def stderr(self) -> str: ...
class CommandChannelMixin:
"""Declares the channel; a driver that can reach a shell on its host implements it."""
if TYPE_CHECKING: # pragma: no cover - declared for type checkers only
def run_command(
self,
command: str,
*,
privileged: bool = False,
timeout: float = 60,
stdin: Optional[bytes] = None,
) -> CommandResult:
"""Run *command* with the host's shell; *stdin*, if given, is written first."""
...
def open_stream(self, command: str, *, privileged: bool = False) -> ByteStream:
"""Start *command* and return a stream to its stdin and stdout."""
...
class ParamikoExecStream:
"""A :class:`ByteStream` over a paramiko exec channel.
stderr is drained on every read: stdout and stderr share one window, and an
unread stderr would stall the stream.
"""
def __init__(self, channel: Any, *, stderr_limit: int = STDERR_TAIL) -> None:
self._channel = channel
self._stderr = b""
self._limit = stderr_limit
def _drain_stderr(self) -> None:
while self._channel.recv_stderr_ready():
self._stderr = (self._stderr + self._channel.recv_stderr(_READ))[-self._limit :]
def read(self, max_bytes: int, timeout: Optional[float] = None) -> bytes:
self._channel.settimeout(timeout)
try:
return self._channel.recv(max_bytes)
except socket.timeout:
raise TimeoutError(f"no data within {timeout}s") from None
finally:
self._drain_stderr()
def write(self, data: bytes) -> None:
self._channel.sendall(data)
def close_write(self) -> None:
self._channel.shutdown_write()
def close(self) -> None:
self._channel.close()
@property
def exit_status(self) -> Optional[int]:
if not self._channel.exit_status_ready():
return None
return self._channel.recv_exit_status()
@property
def stderr(self) -> str:
self._drain_stderr()
return self._stderr.decode("utf-8", "replace")
def _collect(channel: Any, command: str, timeout: float) -> tuple:
out, err = b"", b""
deadline = time.monotonic() + timeout
while True:
if channel.recv_ready():
out += channel.recv(_READ)
elif channel.recv_stderr_ready():
err += channel.recv_stderr(_READ)
elif channel.exit_status_ready():
return out, err
elif time.monotonic() > deadline:
raise TimeoutError(f"{command!r} did not finish within {timeout}s")
else:
time.sleep(_POLL)
def run_on_transport(
transport: Any, command: str, *, stdin: Optional[bytes] = None, timeout: float = 60
) -> CommandResult:
"""Run *command* on an exec channel of the paramiko *transport*."""
channel = transport.open_session()
try:
channel.settimeout(timeout)
channel.exec_command(command)
if stdin is not None:
channel.sendall(stdin)
channel.shutdown_write()
out, err = _collect(channel, command, timeout)
return CommandResult(
out.decode("utf-8", "replace"), err.decode("utf-8", "replace"), channel.recv_exit_status()
)
finally:
channel.close()
def open_stream_on_transport(
transport: Any, command: str, *, stdin_prefix: Optional[bytes] = None
) -> ParamikoExecStream:
"""Start *command* on an exec channel and return its stream.
*stdin_prefix* is written before anything else, which is how a sudo password
reaches ``sudo -S`` ahead of the stream's own bytes.
"""
channel = transport.open_session()
channel.exec_command(command)
if stdin_prefix:
channel.sendall(stdin_prefix)
return ParamikoExecStream(channel)
+141
View File
@@ -0,0 +1,141 @@
# -*- coding: utf-8 -*-
"""Access to a host's container engine, and nothing more.
This package abstracts the connection, the driver builds it, and netOrk does
the talking (NetOrk/netork#765). So there is no container model here and no
parsing:
* :meth:`ContainerEngineMixin.container_engines` says which engines the host
offers, and which API each speaks.
* :meth:`ContainerEngineMixin.open_container_engine` returns a
:class:`ContainerEngineConnection`. It opens a stream to the engine's API
(``docker system dial-stdio``), and it runs the engine's CLI with arguments
the caller chooses. The CLI is there for what the API cannot do: compose,
pulls that need the host user's registry login, and private registry digests.
The driver decides *how* the engine is reached. The hook is
:meth:`ContainerEngineMixin._container_engine_binary`. QNAP returns its
Container Station path, so callers never see where the binary lives. The
command runs through the driver's own :class:`~napalm_device_types.channel.CommandChannelMixin`,
whose privilege handling applies unchanged; Docker normally needs none, because
the login user is in the ``docker`` group.
Mixed in by a driver, not by a role base: a Windows host is an
:class:`~napalm_device_types.os.OSDriver` too, and without a shell channel it
has no ``dial-stdio`` to offer. ``hasattr(driver, "open_container_engine")``
stays truthful.
"""
from __future__ import annotations
import shlex
from typing import Any, Dict, List, Optional, Sequence, TYPE_CHECKING
from napalm_device_types.channel import ByteStream, CommandResult
from napalm_device_types.models import ContainerEngineDict
#: Engines this package knows how to reach, and the API each speaks.
ENGINES: Dict[str, str] = {"docker": "docker-engine"}
_PROBE_TIMEOUT = 15
_CLI_TIMEOUT = 120
class ContainerEngineUnavailable(NotImplementedError):
"""The host offers no such engine, or this package does not know how to reach it."""
class ContainerEngineConnection:
"""The way to one container engine on one host.
Built by :meth:`ContainerEngineMixin.open_container_engine`. Callers pass
arguments; binary, quoting and the channel are the driver's.
"""
def __init__(self, driver: Any, engine: str, binary: str) -> None:
self._driver = driver
self._binary = binary
self.engine = engine
def _command(self, args: Sequence[str]) -> str:
return shlex.join([self._binary, *args])
def open_api(self) -> ByteStream:
"""A stream to the engine's API: HTTP/1.1 over ``system dial-stdio``."""
return self._driver.open_stream(self._command(["system", "dial-stdio"]))
def run_cli(
self,
args: Sequence[str],
*,
stdin: Optional[bytes] = None,
timeout: float = _CLI_TIMEOUT,
privileged: bool = False,
) -> CommandResult:
"""Run the engine's CLI with *args* and wait for it.
*privileged* runs it as root, the way the driver gains root for any
command: for a call the engine refused to the login user.
"""
return self._driver.run_command(
self._command(args), privileged=privileged, timeout=timeout, stdin=stdin
)
def stream_cli(
self, args: Sequence[str], *, merge_stderr: bool = False, privileged: bool = False
) -> ByteStream:
"""Start the engine's CLI with *args* and stream its output.
*merge_stderr* folds stderr into the stream, for tools that report
progress there (``compose up``); *privileged* as for :meth:`run_cli`.
"""
command = self._command(args)
return self._driver.open_stream(
f"{command} 2>&1" if merge_stderr else command, privileged=privileged
)
class ContainerEngineMixin:
"""Adds container engine access to a driver that implements the channel."""
if TYPE_CHECKING: # pragma: no cover - declared for type checkers only
def run_command(
self,
command: str,
*,
privileged: bool = False,
timeout: float = 60,
stdin: Optional[bytes] = None,
) -> CommandResult: ...
def open_stream(self, command: str, *, privileged: bool = False) -> ByteStream: ...
def _container_engine_binary(self, engine: str) -> str:
"""The host command that is *engine*'s CLI. Override where it is not on PATH."""
return engine
def container_engines(self) -> List[ContainerEngineDict]:
"""
Returns the container engines the host offers:
* engine (string) - e.g. "docker"
* api (string) - what its API is, e.g. "docker-engine"
An engine is listed when its CLI exists on the host. Whether the login
user may use it is decided by the first API call, which says
"permission denied" in the stream's stderr.
"""
found: List[ContainerEngineDict] = []
for engine, api in ENGINES.items():
probe = shlex.join(["command", "-v", self._container_engine_binary(engine)])
result = self.run_command(probe, timeout=_PROBE_TIMEOUT)
if result.exit_code == 0 and result.stdout.strip():
found.append({"engine": engine, "api": api})
return found
def open_container_engine(self, engine: str) -> ContainerEngineConnection:
"""The connection to *engine*; :class:`ContainerEngineUnavailable` if unknown."""
if engine not in ENGINES:
raise ContainerEngineUnavailable(f"no way to reach container engine {engine!r}")
return ContainerEngineConnection(self, engine, self._container_engine_binary(engine))
+14 -5
View File
@@ -14,6 +14,7 @@ from typing import Any, Dict, List, TYPE_CHECKING
from napalm_device_types.base import DeviceTypeDriver
from napalm_device_types.packages import PackageManagementMixin
from napalm_device_types.health_metrics import HealthMetricsMixin
from napalm_device_types.provisioning import QEMU_GUEST_AGENTS, GuestAgents
from napalm_device_types.models import (
NICConfigDict,
NetworkTargetDict,
@@ -44,11 +45,11 @@ class HypervisorDriver(PackageManagementMixin, HealthMetricsMixin, DeviceTypeDri
ROLE: str = "hypervisor"
TYPE_LABEL: str = "Hypervisor"
#: What cloud-init installs and starts on a VM provisioned through this
#: driver, so the hypervisor can read the guest's IP address back.
GUEST_AGENT_PACKAGES: tuple[str, ...] = ("qemu-guest-agent",)
GUEST_AGENT_RUNCMD: tuple[str, ...] = ("systemctl enable --now qemu-guest-agent",)
#: What cloud-init installs and runs on a VM provisioned through this
#: driver, per guest OS, so the hypervisor can read the guest's IP
#: address back. The guest operating systems listed here are the ones
#: the driver can provision; netOrk reads this off the class.
GUEST_AGENTS: GuestAgents = {"linux": QEMU_GUEST_AGENTS["linux"]}
# ------------------------------------------------------------------
@@ -464,6 +465,7 @@ class HypervisorDriver(PackageManagementMixin, HealthMetricsMixin, DeviceTypeDri
disk_resize_gb: int | None = None,
storage: str | None = None,
cpu_type: str | None = None,
guest_os: str = "linux",
download_timeout: int = 300,
timeout: int = 180,
) -> VMProvisionResultDict:
@@ -512,6 +514,13 @@ class HypervisorDriver(PackageManagementMixin, HealthMetricsMixin, DeviceTypeDri
value other than None with ValueError; one that has it raises
ValueError for a name it does not list or that is not
``available`` on this node, before creating anything.
guest_os (string) - the operating system in the image, a key of
``GUEST_AGENTS`` (``"linux"``, ``"freebsd"``, ``"openbsd"``). The
driver sets up the VM's hardware for it (OS type, how the guest
agent is attached) and, where the guest needs it, writes the
network-config itself (``provisioning.network_config``). A
guest_os the driver does not list raises ValueError before
anything is created. Default ``"linux"``.
download_timeout (int) - maximum seconds to wait for the image download
(skipped entirely if already cached on the hypervisor). Default 300.
timeout (int) - maximum seconds to wait for the remaining provisioning
+276
View File
@@ -0,0 +1,276 @@
# -*- coding: utf-8 -*-
"""Listening sockets: what listens on which address, and which service it is.
Whether a service is reachable from outside its host is decided by what it
listens on -- ``0.0.0.0:5432`` is, ``127.0.0.1:5432`` is not -- and that is read
the same way on every Linux host. So the command and its parse live here once,
and a driver only carries the command across.
**One round trip.** ``ss -lntup`` lists every listening TCP and bound UDP
socket with the processes holding it; for each of those processes,
``/proc/<pid>/cgroup`` says which systemd service or container it runs in. The
report is framed, and a report whose end is missing raises: a list cut short
must never read as sockets that closed.
**Root, and without it.** Only root sees every process behind a socket.
The whole script therefore goes to the host as one ``sh -c`` argument -- a
driver that prefixes ``sudo -n`` would otherwise run only its first command as
root. When that call brings no report back (no sudo, a wrong password), the
command runs again without privilege: the sockets are still worth having, and
the reading says it is not ``attributed``.
**No ``-H``.** iproute2 before 4.10 has no option to leave out the header and
fails on it, which would read as nothing listening. The parse skips the header
instead.
**Without ``ss``, ``netstat``.** OpenWrt's busybox and old net-tools hosts have
no ``ss``; ``netstat -lntup`` lists the same sockets with ``PID/Program``, and
the cgroups are read for its PIDs alike. On OpenWrt the cgroup names the procd
service (``/services/<name>/<instance>``), a jailed one too, whose PID is not the
one procd reports. A host with neither raises :class:`ListeningSocketsUnavailable`.
"""
from __future__ import annotations
import re
from shlex import quote
from typing import Dict, List, Optional, Tuple, TYPE_CHECKING
from napalm_device_types.models import ListeningSocketDict, ListeningSocketsDict
from napalm_device_types.terminal import strip_terminal_codes
_BEGIN = "SOCK_BEGIN"
_END = "SOCK_END"
_NO_SS = "no-ss"
_RC_RE = re.compile(r"^__SS_RC=(\d+)$")
#: One line, POSIX ``sh``, read-only. The frame markers are printed in two
#: halves so that a transport which echoes the command does not show them early.
#: ``ss`` is in ``/usr/sbin`` on some systems, outside a login user's ``PATH``.
#: ``netstat`` names a socket's process as ``PID/Program``, ``ss`` as ``pid=PID``.
LISTENING_SOCKETS_COMMAND = (
"PATH=$PATH:/usr/sbin:/sbin; "
"printf '%s%s\\n' SOCK_ BEGIN; "
"if command -v ss >/dev/null 2>&1; then "
"t=ss; s=$(ss -lntup 2>&1); r=$?; "
"pids=$(printf '%s\\n' \"$s\" | grep -o 'pid=[0-9]*' | cut -d= -f2); "
"elif command -v netstat >/dev/null 2>&1; then "
"t=netstat; s=$(netstat -lntup 2>&1); r=$?; "
"pids=$(printf '%s\\n' \"$s\" | grep -o ' [0-9][0-9]*/' | tr -d ' /'); "
"else t=; fi; "
"if [ -n \"$t\" ]; then "
"echo \"[$t]\"; printf '%s\\n' \"$s\"; echo \"__SS_RC=$r\"; "
"echo '[cgroups]'; "
"for p in $(printf '%s\\n' \"$pids\" | sort -u); do "
"sed \"s|^|$p |\" /proc/$p/cgroup 2>/dev/null; done; "
"else echo '[no-ss]'; fi; "
"printf '%s%s\\n' SOCK_ END"
)
_PROTOCOLS = frozenset({"tcp", "udp"})
#: netstat names the IPv6 sockets of net-tools ``tcp6``/``udp6``; busybox does not.
_NETSTAT_PROTOCOLS = {"tcp": "tcp", "tcp6": "tcp", "udp": "udp", "udp6": "udp"}
#: ``1604/dropbear``; net-tools prints ``700/sshd: /usr/sbin``.
_PROGRAM_RE = re.compile(r"^(\d+)/(\S*)")
#: ``users:(("nginx",pid=901,fd=6),("nginx",pid=900,fd=6))``
_USER_RE = re.compile(r'\("((?:[^"\\]|\\.)*)",pid=(\d+),fd=\d+\)')
#: The service a cgroup path runs in: its deepest ``*.service`` component.
_SERVICE_RE = re.compile(r"/([^/]+)\.service(?=/|$)")
#: OpenWrt's procd: ``/services/<name>/<instance>``.
_PROCD_RE = re.compile(r"^/services/([^/]+)(?:/|$)")
#: A container's cgroup: ``docker-<id>.scope`` (systemd driver), ``/docker/<id>`` (cgroupfs).
_CONTAINER_RE = re.compile(
r"(?:docker|libpod)-([0-9a-f]{64})\.scope|/(?:docker|libpod)/([0-9a-f]{64})(?=/|$)"
)
class ListeningSocketsUnavailable(NotImplementedError):
"""The host has neither ``ss`` nor ``netstat``; nothing to read, nothing to retry."""
def _frame(output: str) -> List[str]:
lines = [line.strip() for line in strip_terminal_codes(output).splitlines()]
try:
start = lines.index(_BEGIN)
end = lines.index(_END, start)
except ValueError:
raise ValueError("no intact listening socket report in the output") from None
return lines[start + 1 : end]
def _sections(lines: List[str]) -> Dict[str, List[str]]:
sections: Dict[str, List[str]] = {}
current: List[str] = []
for line in lines:
if line.startswith("[") and line.endswith("]") and " " not in line:
current = sections.setdefault(line[1:-1], [])
else:
current.append(line)
return sections
def _split_local(local: str) -> Optional[Tuple[str, Optional[str], int]]:
"""``[fe80::1%eth0]:546`` -> ``("fe80::1", "eth0", 546)``; None if no port."""
host, sep, port = local.rpartition(":")
if not sep or not port.isdigit():
return None
if host.startswith("[") and host.endswith("]"):
host = host[1:-1]
address, _, zone = host.partition("%")
return address or "*", zone or None, int(port)
def _cgroup_paths(lines: List[str]) -> Dict[int, str]:
"""Each process's cgroup path: the unified hierarchy, or systemd's under v1."""
paths: Dict[int, str] = {}
for line in lines:
pid, _, entry = line.partition(" ")
parts = entry.split(":", 2)
if not pid.isdigit() or len(parts) != 3:
continue
hierarchy, controllers, path = parts
if (hierarchy == "0" and controllers == "") or controllers == "name=systemd":
paths[int(pid)] = path
return paths
def _unit(path: Optional[str]) -> Optional[str]:
services = _SERVICE_RE.findall(path or "")
if services:
return str(services[-1])
procd = _PROCD_RE.match(path or "")
return str(procd.group(1)) if procd else None
def _container(path: Optional[str]) -> Optional[str]:
match = _CONTAINER_RE.search(path or "")
return (match.group(1) or match.group(2)) if match else None
def _entry(
proto: str,
local: Tuple[str, Optional[str], int],
process: Optional[str],
pid: Optional[int],
paths: Dict[int, str],
) -> ListeningSocketDict:
address, interface, port = local
path = paths.get(pid) if pid is not None else None
return {
"proto": proto,
"address": address,
"port": port,
"interface": interface,
"process": process,
"pid": pid,
"unit": _unit(path),
"container_id": _container(path),
}
def _ss_socket(line: str, paths: Dict[int, str]) -> Optional[ListeningSocketDict]:
parts = line.split()
if len(parts) < 5 or parts[0] not in _PROTOCOLS:
return None
local = _split_local(parts[4])
if local is None:
return None
users = _USER_RE.findall(line)
process, pid = (users[0][0], int(users[0][1])) if users else (None, None)
return _entry(parts[0], local, process, pid, paths)
def _netstat_socket(line: str, paths: Dict[int, str]) -> Optional[ListeningSocketDict]:
"""``Proto Recv-Q Send-Q Local Foreign [State] PID/Program`` -- a UDP line
has no state, and ``-`` is a socket without a process."""
parts = line.split()
proto = _NETSTAT_PROTOCOLS.get(parts[0]) if parts else None
if proto is None or len(parts) < 6:
return None
local = _split_local(parts[3])
if local is None:
return None
process, pid = None, None
for token in parts[5:7]:
program = _PROGRAM_RE.match(token)
if program:
pid, process = int(program.group(1)), program.group(2).rstrip(":") or None
break
return _entry(proto, local, process, pid, paths)
_PARSERS = {"ss": _ss_socket, "netstat": _netstat_socket}
def parse_listening_sockets(output: str) -> List[ListeningSocketDict]:
"""Parse what :data:`LISTENING_SOCKETS_COMMAND` printed, sorted by protocol,
port and address.
:raises ListeningSocketsUnavailable: when the host has neither ``ss`` nor ``netstat``.
:raises ValueError: when the output carries no intact report, or the tool failed.
"""
sections = _sections(_frame(output))
if _NO_SS in sections:
raise ListeningSocketsUnavailable("the host has neither ss nor netstat")
tool = "netstat" if "netstat" in sections else "ss"
lines = sections.get(tool, [])
statuses = [m.group(1) for m in map(_RC_RE.match, lines) if m]
if not statuses or statuses[-1] != "0":
detail = " ".join(line for line in lines if not _RC_RE.match(line))[:200]
raise ValueError(f"{tool} did not list the sockets: {detail or 'no exit status'}")
paths = _cgroup_paths(sections.get("cgroups", []))
parse = _PARSERS[tool]
sockets = [s for s in (parse(line, paths) for line in lines) if s is not None]
return sorted(sockets, key=lambda s: (s["proto"], s["port"], s["address"], s["interface"] or ""))
class ListeningSocketsMixin:
"""Adds :meth:`get_listening_sockets` to a driver that can run a command on a Linux host.
With ``ss``, or ``netstat`` where there is none (OpenWrt's busybox).
The template form (README, "Function classes"): the command and its parse
are the same everywhere, so they are concrete here, and a driver supplies
only :meth:`_run_listening_sockets_command` -- how a command reaches its
host, and how it gains root there. Mixed in by the drivers that can, not by
:class:`~napalm_device_types.os.OSDriver`: a Windows host is an OS driver
too, and ``hasattr(driver, "get_listening_sockets")`` has to stay truthful.
"""
if TYPE_CHECKING: # pragma: no cover - declared for type checkers only
def _run_listening_sockets_command(self, command: str, *, privileged: bool) -> str:
"""Run *command* on the host and return what it printed; as root
when *privileged*. The command is a single ``sh -c`` invocation."""
...
def get_listening_sockets(self) -> ListeningSocketsDict:
"""
Returns every listening TCP and bound UDP socket, with the process,
systemd service and container behind it.
* attributed (bool) - read as root, so every process is named
* sockets (list) - see :class:`~napalm_device_types.models.ListeningSocketDict`
Example::
{
"attributed": True,
"sockets": [
{"proto": "tcp", "address": "0.0.0.0", "port": 5432,
"interface": None, "process": "postgres", "pid": 812,
"unit": "postgresql@16-main", "container_id": None},
],
}
:raises ListeningSocketsUnavailable: if the host has neither ``ss`` nor ``netstat``.
:raises ValueError: if neither reading carried an intact report.
"""
command = f"sh -c {quote(LISTENING_SOCKETS_COMMAND)}"
try:
output = self._run_listening_sockets_command(command, privileged=True)
return {"attributed": True, "sockets": parse_listening_sockets(output)}
except ValueError:
pass
output = self._run_listening_sockets_command(command, privileged=False)
return {"attributed": False, "sockets": parse_listening_sockets(output)}
+42
View File
@@ -335,6 +335,37 @@ class KernelFactsDict(TypedDict):
config: Optional[Dict[str, str]] # build configuration, set options only; quotes stripped
class ListeningSocketDict(TypedDict):
"""A TCP socket that listens, or a UDP socket that is bound, on the host.
One entry per socket as ``ss`` lists it. ``address`` is printed the way ss
prints it, without brackets or zone: ``0.0.0.0`` and ``::`` are every
address of their family, ``*`` every address of both. Whether that is
reachable from outside the host is the consumer's call.
"""
proto: str # "tcp" or "udp"
address: str # "0.0.0.0", "::", "*", "127.0.0.1", "::ffff:127.0.0.1", ...
port: int
interface: Optional[str] # the %zone a socket is bound to ("lo", "eth0"), if any
process: Optional[str] # the first process holding it; None without one or without root
pid: Optional[int]
unit: Optional[str] # the process's systemd service, without ".service"
container_id: Optional[str] # the full container ID, for a container on the host network
class ListeningSocketsDict(TypedDict):
"""What ``ListeningSocketsMixin.get_listening_sockets`` read.
``attributed`` is False when the reading ran without root: ``ss`` then names
only the login user's own processes, so a socket without a process means
"not told", not "the kernel's".
"""
attributed: bool
sockets: List[ListeningSocketDict]
class PortForwardDict(TypedDict):
"""A port the WAN side can reach, forwarded to a host inside.
@@ -840,6 +871,17 @@ class DockerInfoDict(TypedDict):
outdated_images: NotRequired[List[str]] # image names with a newer remote digest
class ContainerEngineDict(TypedDict):
"""One container engine a host offers (``container_engines()``).
Only how to reach it: which engine, and which API it speaks. What runs on
it is read and modelled above the driver, in netOrk (netork#765).
"""
engine: str # "docker"
api: str # "docker-engine": the Docker Engine API over ``system dial-stdio``
class DeviceActionResultDict(TypedDict):
"""Return value of ``run_device_action()``."""
+64
View File
@@ -0,0 +1,64 @@
"""
What provisioning a VM from a cloud image needs to know about the guest.
The same for every hypervisor: which agent a guest runs so a QEMU-based
hypervisor can read its IP, the network-config cloud-init and FreeBSD's
nuageinit both read, and how a packed cloud image is unpacked. How a
hypervisor attaches, boots and talks to the VM stays in its driver.
"""
from __future__ import annotations
from collections.abc import Mapping
from typing import Any
#: ``{guest_os: (packages, runcmd)}``: what cloud-init installs and runs so
#: the hypervisor can read the guest's IP address back.
GuestAgents = Mapping[str, tuple[tuple[str, ...], tuple[str, ...]]]
#: The QEMU guest agent per guest OS, for hypervisors built on QEMU. Package
#: names and service commands were verified on FreeBSD 15.1 and OpenBSD 7.9
#: cloud images (NetOrk/netork#793).
QEMU_GUEST_AGENTS: GuestAgents = {
"linux": (("qemu-guest-agent",), ("systemctl enable --now qemu-guest-agent",)),
"freebsd": (
("qemu-guest-agent",),
("sysrc qemu_guest_agent_enable=YES", "service qemu-guest-agent start"),
),
"openbsd": (("qemu-ga",), ("rcctl enable qemu_ga", "rcctl start qemu_ga")),
}
#: Suffix of a packed image -> the command that writes it unpacked to stdout.
_DECOMPRESSORS = {
"xz": "xz -dc",
"gz": "gzip -dc",
"bz2": "bzip2 -dc",
"zst": "zstd -dc",
}
def network_config(nics: list[tuple[str, bool]]) -> dict[str, Any] | None:
"""Network-config v2: DHCP on every ``(mac, dhcp)`` NIC that asks for it.
Version 2 because FreeBSD's nuageinit reads no other: given the v1 that
Proxmox generates, it fails and skips the rest of its first stage,
runcmd included. ``None`` when no NIC wants DHCP, so the guest is left
alone rather than told to configure nothing.
"""
ethernets = {
f"nic{i}": {"match": {"macaddress": mac.lower()}, "dhcp4": True}
for i, (mac, dhcp) in enumerate(nics)
if dhcp
}
return {"version": 2, "ethernets": ethernets} if ethernets else None
def split_compression(filename: str) -> tuple[str, str | None]:
"""``(unpacked filename, command that unpacks to stdout)`` for *filename*.
The command is None for an image that is not packed, which is then
imported as it is.
"""
stem, dot, suffix = filename.rpartition(".")
command = _DECOMPRESSORS.get(suffix.lower()) if dot and stem else None
return (stem, command) if command else (filename, None)
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "napalm-device-types"
version = "2.3.0"
version = "3.0.0"
description = "Abstract device-type base classes for NAPALM drivers"
readme = "README.md"
requires-python = ">=3.10"
+169
View File
@@ -0,0 +1,169 @@
"""The command channel: how a driver runs a command and opens a byte stream.
netOrk's container runtime driver speaks the Docker Engine API over a stream
from ``docker system dial-stdio`` (NetOrk/netork#765). The driver only provides
the way there; these are the SSH pieces every SSH driver can reuse. The
paramiko channel is faked, so the tests pin behaviour, not paramiko.
"""
from __future__ import annotations
import socket
import pytest
from napalm_device_types.channel import (
CommandResult,
ParamikoExecStream,
open_stream_on_transport,
run_on_transport,
)
class FakeChannel:
"""Just enough of paramiko.Channel: scripted stdout/stderr chunks, an exit code."""
def __init__(self, stdout=(), stderr=(), exit_code=0, hang=False):
self.out = list(stdout)
self.err = list(stderr)
self.exit_code = exit_code
self.hang = hang
self.sent = b""
self.command = None
self.write_closed = False
self.closed = False
self.timeout = None
def exec_command(self, command):
self.command = command
def settimeout(self, timeout):
self.timeout = timeout
def sendall(self, data):
self.sent += data
def shutdown_write(self):
self.write_closed = True
def close(self):
self.closed = True
def recv_ready(self):
return bool(self.out)
def recv(self, n):
if self.out:
return self.out.pop(0)
if self.hang:
raise socket.timeout()
return b""
def recv_stderr_ready(self):
return bool(self.err)
def recv_stderr(self, n):
return self.err.pop(0) if self.err else b""
def exit_status_ready(self):
return not self.hang and not self.out and not self.err
def recv_exit_status(self):
return self.exit_code
class FakeTransport:
def __init__(self, channel):
self.channel = channel
def open_session(self):
return self.channel
def test_run_collects_stdout_stderr_and_the_exit_code():
ch = FakeChannel(stdout=[b"hel", b"lo\n"], stderr=[b"warn\n"], exit_code=3)
result = run_on_transport(FakeTransport(ch), "echo hello", timeout=5)
assert result == CommandResult(stdout="hello\n", stderr="warn\n", exit_code=3)
assert ch.command == "echo hello"
assert ch.closed
def test_run_sends_stdin_and_then_closes_the_write_side():
ch = FakeChannel(stdout=[b"ok"])
run_on_transport(FakeTransport(ch), "sudo -S true", stdin=b"pw\n", timeout=5)
assert ch.sent == b"pw\n"
assert ch.write_closed
def test_run_raises_when_the_command_does_not_finish_in_time():
ch = FakeChannel(hang=True)
with pytest.raises(TimeoutError, match="did not finish"):
run_on_transport(FakeTransport(ch), "sleep 999", timeout=0.05)
assert ch.closed
def test_a_stream_reads_until_eof_and_drains_stderr_on_the_way():
ch = FakeChannel(stdout=[b"HTTP/1.1 200 OK\r\n"], stderr=[b"note\n"])
stream = open_stream_on_transport(FakeTransport(ch), "docker system dial-stdio")
assert stream.read(4096, timeout=5) == b"HTTP/1.1 200 OK\r\n"
assert stream.read(4096, timeout=5) == b""
assert stream.stderr == "note\n"
assert ch.command == "docker system dial-stdio"
def test_a_stream_writes_and_closes_its_write_side_only_when_asked():
"""dial-stdio answers HTTP 499 to a request whose writer closed early
(netork#771), so write() never implies close_write()."""
ch = FakeChannel()
stream = open_stream_on_transport(FakeTransport(ch), "docker system dial-stdio")
stream.write(b"GET /_ping HTTP/1.1\r\n\r\n")
assert not ch.write_closed
stream.close_write()
assert ch.write_closed and ch.sent.startswith(b"GET /_ping")
def test_a_stream_sends_its_stdin_prefix_first():
"""How a sudo password reaches `sudo -S` before the stream's own bytes."""
ch = FakeChannel()
open_stream_on_transport(FakeTransport(ch), "sudo -S -p '' cmd", stdin_prefix=b"pw\n")
assert ch.sent == b"pw\n"
def test_a_stream_times_out_as_timeout_error():
stream = ParamikoExecStream(FakeChannel(hang=True))
with pytest.raises(TimeoutError):
stream.read(10, timeout=0.01)
def test_the_stream_keeps_only_a_bounded_tail_of_stderr():
ch = FakeChannel(stderr=[b"x" * 100, b"y" * 100])
stream = ParamikoExecStream(ch, stderr_limit=50)
assert stream.stderr == "y" * 50
def test_exit_status_is_none_while_running_and_the_code_after():
running = ParamikoExecStream(FakeChannel(hang=True))
finished = ParamikoExecStream(FakeChannel(exit_code=1))
assert running.exit_status is None
assert finished.exit_status == 1
def test_close_closes_the_channel():
ch = FakeChannel()
ParamikoExecStream(ch).close()
assert ch.closed
+159
View File
@@ -0,0 +1,159 @@
"""Access to a container engine: which engines a host has, and a way to each.
This package abstracts the connection, the driver builds it, and netOrk does
the talking (NetOrk/netork#765, decided 2026-10-07). So nothing here knows what
a container is. It knows how to reach the engine's API (``dial-stdio``) and how
to run the engine's CLI with arguments netOrk chooses, and the driver alone
knows where that CLI lives.
"""
from __future__ import annotations
import pytest
from napalm_device_types import OSDriver
from napalm_device_types.channel import CommandResult
from napalm_device_types.container_engine import (
ContainerEngineConnection,
ContainerEngineMixin,
ContainerEngineUnavailable,
)
class FakeDriver(ContainerEngineMixin):
"""A driver with a channel: records what it is asked to run."""
def __init__(self, present=("docker",), binary=None):
self.present = set(present)
self.binary = binary
self.commands = []
self.streams = []
def run_command(self, command, *, privileged=False, timeout=60, stdin=None):
self.commands.append((command, privileged, timeout, stdin))
found = any(command == f"command -v {b}" for b in self._binaries())
return CommandResult(stdout="/usr/bin/x\n" if found else "", stderr="", exit_code=0 if found else 1)
def open_stream(self, command, *, privileged=False):
self.streams.append((command, privileged))
return object()
def _binaries(self):
return {self._container_engine_binary(e) for e in self.present}
def _container_engine_binary(self, engine):
return self.binary or super()._container_engine_binary(engine)
def test_a_host_with_docker_lists_it():
assert FakeDriver().container_engines() == [{"engine": "docker", "api": "docker-engine"}]
def test_a_host_without_the_cli_lists_nothing():
assert FakeDriver(present=()).container_engines() == []
def test_the_api_is_reached_through_dial_stdio():
driver = FakeDriver()
driver.open_container_engine("docker").open_api()
assert driver.streams == [("docker system dial-stdio", False)]
def test_the_driver_decides_where_the_binary_lives():
"""QNAP keeps docker under Container Station's path; netOrk never sees it."""
path = "/share/CACHEDEV1_DATA/.qpkg/container-station/bin/docker"
driver = FakeDriver(binary=path)
assert driver.container_engines() == [{"engine": "docker", "api": "docker-engine"}]
driver.open_container_engine("docker").open_api()
assert driver.streams == [(f"{path} system dial-stdio", False)]
def test_cli_arguments_are_quoted_and_prefixed_by_the_binary():
driver = FakeDriver()
conn = driver.open_container_engine("docker")
conn.run_cli(["compose", "-f", "/srv/my stack/compose.yml", "config", "--format", "json"])
command, privileged, timeout, stdin = driver.commands[-1]
assert command == "docker compose -f '/srv/my stack/compose.yml' config --format json"
assert privileged is False and stdin is None
def test_cli_stdin_and_timeout_reach_the_channel():
driver = FakeDriver()
driver.open_container_engine("docker").run_cli(["login", "--password-stdin"], stdin=b"x", timeout=30)
_command, _privileged, timeout, stdin = driver.commands[-1]
assert (timeout, stdin) == (30, b"x")
def test_a_streamed_cli_call_can_merge_stderr():
"""compose writes its progress to stderr; a merged stream carries both."""
driver = FakeDriver()
conn = driver.open_container_engine("docker")
conn.stream_cli(["pull", "redis:7-alpine"])
conn.stream_cli(["compose", "up", "-d", "db"], merge_stderr=True)
assert driver.streams == [
("docker pull redis:7-alpine", False),
("docker compose up -d db 2>&1", False),
]
def test_a_shell_metacharacter_in_an_argument_stays_an_argument():
driver = FakeDriver()
driver.open_container_engine("docker").run_cli(["inspect", "x; rm -rf /"])
assert driver.commands[-1][0] == "docker inspect 'x; rm -rf /'"
def test_an_unknown_engine_is_refused():
with pytest.raises(ContainerEngineUnavailable):
FakeDriver().open_container_engine("rkt")
def test_the_connection_names_its_engine():
conn = FakeDriver().open_container_engine("docker")
assert isinstance(conn, ContainerEngineConnection)
assert conn.engine == "docker"
def test_an_os_driver_does_not_get_container_engine_access_by_role():
"""Mixed in by the driver, not by OSDriver: a Windows host is an OS driver
too, and `hasattr(driver, "open_container_engine")` has to stay truthful."""
assert not issubclass(OSDriver, ContainerEngineMixin)
assert not hasattr(OSDriver, "open_container_engine")
class TestPrivilegedCli:
"""A refused engine call can be repeated as root (netork#773): the driver
knows the binary and how it gains root, so the connection carries it."""
def test_run_cli_passes_privileged_to_the_channel(self):
driver = FakeDriver()
driver.open_container_engine("docker").run_cli(["restart", "web"], privileged=True)
assert driver.commands == [("docker restart web", True, 120, None)]
def test_run_cli_is_unprivileged_by_default(self):
driver = FakeDriver()
driver.open_container_engine("docker").run_cli(["ps"])
assert driver.commands[0][1] is False
def test_stream_cli_passes_privileged_too(self):
driver = FakeDriver()
driver.open_container_engine("docker").stream_cli(["pull", "nginx"], privileged=True)
assert driver.streams == [("docker pull nginx", True)]
+396
View File
@@ -0,0 +1,396 @@
"""get_listening_sockets: what listens on which address, and which service it is.
Whether a service is reachable from outside its host is decided by what it
listens on -- ``0.0.0.0:5432`` is, ``127.0.0.1:5432`` is not. Reading that is
the same on every Linux host: ``ss`` for the sockets, ``/proc/<pid>/cgroup``
for the systemd unit or container a process belongs to. So both live here once,
and a driver only carries the command across (#658 in netOrk).
A host without ``ss`` -- OpenWrt's busybox, an old net-tools box -- is read
with ``netstat -lntup`` instead, and on OpenWrt the cgroup names the procd
service (#673 in netOrk).
"""
from __future__ import annotations
import os
import shutil
import subprocess
import pytest
from napalm_device_types import OSDriver
from napalm_device_types.listening import (
LISTENING_SOCKETS_COMMAND,
ListeningSocketsMixin,
ListeningSocketsUnavailable,
parse_listening_sockets,
)
CID = "4f1c" + "0" * 60
SS = """Netid State Recv-Q Send-Q Local Address:Port Peer Address:PortProcess
udp UNCONN 0 0 127.0.0.53%lo:53 0.0.0.0:* users:(("systemd-resolve",pid=612,fd=13))
udp UNCONN 0 0 0.0.0.0%ens18:68 0.0.0.0:* users:(("systemd-network",pid=590,fd=22))
tcp LISTEN 0 4096 0.0.0.0:5432 0.0.0.0:* users:(("postgres",pid=812,fd=6))
tcp LISTEN 0 128 [::]:22 [::]:* users:(("sshd",pid=700,fd=4))
tcp LISTEN 0 511 *:80 *:* users:(("nginx",pid=901,fd=6),("nginx",pid=900,fd=6))
tcp LISTEN 0 4096 [::ffff:127.0.0.1]:8125 *:* users:(("statsd",pid=950,fd=3))
tcp LISTEN 0 4096 0.0.0.0:8080 0.0.0.0:* users:(("docker-proxy",pid=1200,fd=4))
tcp LISTEN 0 4096 0.0.0.0:9100 0.0.0.0:* users:(("node_exporter",pid=1300,fd=3))
tcp LISTEN 0 64 0.0.0.0:2049 0.0.0.0:*
"""
CGROUPS = f"""612 0::/system.slice/systemd-resolved.service
590 0::/system.slice/systemd-networkd.service
812 0::/system.slice/system-postgresql.slice/postgresql@16-main.service
700 0::/system.slice/ssh.service
900 0::/system.slice/nginx.service
901 0::/system.slice/nginx.service
950 0::/user.slice/user-1000.slice/session-3.scope
1200 0::/system.slice/docker.service
1300 0::/system.slice/docker-{CID}.scope
"""
def _wire(ss: str = SS, cgroups: str = CGROUPS, *, rc: int = 0, noise: str = "") -> str:
"""The report as the command prints it, framed."""
return f"{noise}SOCK_BEGIN\n[ss]\n{ss}__SS_RC={rc}\n[cgroups]\n{cgroups}SOCK_END\n"
def _by_port(sockets: list) -> dict:
return {(s["proto"], s["port"], s["address"]): s for s in sockets}
class TestParsing:
def test_a_socket_comes_with_its_process_and_unit(self):
sockets = _by_port(parse_listening_sockets(_wire()))
assert sockets[("tcp", 5432, "0.0.0.0")] == {
"proto": "tcp",
"address": "0.0.0.0",
"port": 5432,
"interface": None,
"process": "postgres",
"pid": 812,
"unit": "postgresql@16-main",
"container_id": None,
}
@pytest.mark.parametrize(
"local, address, interface",
[
("[::]:22", "::", None),
("*:80", "*", None),
("127.0.0.53%lo:53", "127.0.0.53", "lo"),
("0.0.0.0%ens18:68", "0.0.0.0", "ens18"),
("[::ffff:127.0.0.1]:8125", "::ffff:127.0.0.1", None),
("[fe80::1%eth0]:546", "fe80::1", "eth0"),
(":::22", "::", None), # iproute2 4.9 prints IPv6 without brackets
],
)
def test_every_address_form(self, local: str, address: str, interface):
line = f"tcp LISTEN 0 128 {local} *:* users:((\"x\",pid=1,fd=3))\n"
[socket] = parse_listening_sockets(_wire(line, "1 0::/system.slice/x.service\n"))
assert (socket["address"], socket["interface"]) == (address, interface)
def test_the_first_of_several_processes_names_the_socket(self):
nginx = _by_port(parse_listening_sockets(_wire()))[("tcp", 80, "*")]
assert (nginx["process"], nginx["pid"], nginx["unit"]) == ("nginx", 901, "nginx")
def test_a_socket_without_a_process_is_kept(self):
"""The kernel's nfsd, or every socket of another user without root."""
nfs = _by_port(parse_listening_sockets(_wire()))[("tcp", 2049, "0.0.0.0")]
assert (nfs["process"], nfs["pid"], nfs["unit"], nfs["container_id"]) == (
None,
None,
None,
None,
)
def test_the_header_is_no_socket(self):
sockets = parse_listening_sockets(_wire())
assert len(sockets) == 9
assert all(s["proto"] in ("tcp", "udp") for s in sockets)
def test_sorted_by_protocol_port_and_address(self):
keys = [(s["proto"], s["port"], s["address"]) for s in parse_listening_sockets(_wire())]
assert keys == sorted(keys)
def test_nothing_listening_is_an_empty_list(self):
assert parse_listening_sockets(_wire(SS.splitlines(True)[0], "")) == []
class TestCgroups:
def _unit_and_container(self, cgroup_lines: str) -> tuple:
line = 'tcp LISTEN 0 128 0.0.0.0:1 0.0.0.0:* users:(("p",pid=5,fd=3))\n'
[socket] = parse_listening_sockets(_wire(line, cgroup_lines))
return socket["unit"], socket["container_id"]
def test_a_container_on_the_host_network(self):
assert self._unit_and_container(f"5 0::/system.slice/docker-{CID}.scope\n") == (
None,
CID,
)
def test_a_container_under_the_cgroupfs_driver(self):
assert self._unit_and_container(f"5 0::/docker/{CID}\n") == (None, CID)
def test_cgroup_v1_reads_the_systemd_hierarchy(self):
lines = "5 12:pids:/user.slice\n5 1:name=systemd:/system.slice/ssh.service\n"
assert self._unit_and_container(lines) == ("ssh", None)
def test_an_escaped_template_instance(self):
lines = "5 0::/system.slice/system-wg\\x2dquick.slice/wg-quick@wg0.service\n"
assert self._unit_and_container(lines) == ("wg-quick@wg0", None)
def test_a_procd_service_on_openwrt(self):
"""procd puts every instance into /services/<name>/<instance>, a jailed
one too -- its PID is not the one procd reports, its cgroup is."""
assert self._unit_and_container("5 0::/services/dnsmasq/cfg01411c\n") == (
"dnsmasq",
None,
)
def test_a_login_session_is_no_unit(self):
assert self._unit_and_container("5 0::/user.slice/user-1000.slice/session-3.scope\n") == (
None,
None,
)
def test_a_process_gone_before_its_cgroup_was_read(self):
assert self._unit_and_container("") == (None, None)
def test_the_docker_daemons_own_proxy_stays_raw(self):
"""docker-proxy runs in docker.service. Telling it apart from a real
listener is the consumer's business; the reading reports what is there."""
proxy = _by_port(parse_listening_sockets(_wire()))[("tcp", 8080, "0.0.0.0")]
assert (proxy["process"], proxy["unit"]) == ("docker-proxy", "docker")
class TestFailures:
def test_whatever_surrounds_the_frame_is_ignored(self):
noisy = _wire(noise="$ sh -c '...'\nWelcome to Ubuntu\n")
assert len(parse_listening_sockets(noisy)) == 9
def test_output_without_the_frame_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets("sudo: a password is required\n")
def test_a_report_cut_short_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets(_wire()[:-len("SOCK_END\n")])
def test_a_host_without_ss_or_netstat_is_unavailable(self):
with pytest.raises(ListeningSocketsUnavailable):
parse_listening_sockets("SOCK_BEGIN\n[no-ss]\nSOCK_END\n")
def test_ss_failing_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets(_wire("ss: invalid option -- 'p'\n", "", rc=1))
def test_unavailable_is_not_a_value_error(self):
"""A caller retries a broken report without privilege; a host without
ss has nothing to retry."""
assert not issubclass(ListeningSocketsUnavailable, ValueError)
# busybox netstat on OpenWrt 25.12, addresses replaced by documentation ones.
BUSYBOX = """Active Internet connections (only servers)
Proto Recv-Q Send-Q Local Address Foreign Address State PID/Program name
tcp 0 0 0.0.0.0:22 0.0.0.0:* LISTEN 1604/dropbear
tcp 0 0 0.0.0.0:443 0.0.0.0:* LISTEN 1969/uhttpd
tcp 0 0 192.0.2.15:53 0.0.0.0:* LISTEN 1495/dnsmasq
tcp 0 0 127.0.0.1:53 0.0.0.0:* LISTEN 1495/dnsmasq
tcp 0 0 :::22 :::* LISTEN 1604/dropbear
tcp 0 0 fe80::1:53 :::* LISTEN 1495/dnsmasq
tcp 0 0 ::1:53 :::* LISTEN 1495/dnsmasq
udp 0 0 192.0.2.15:53 0.0.0.0:* 1495/dnsmasq
udp 0 0 0.0.0.0:161 0.0.0.0:* 3173/snmpd
udp 0 0 0.0.0.0:5353 0.0.0.0:* -
"""
PROCD = """1495 0::/services/dnsmasq/cfg01411c
1604 0::/services/dropbear/instance1
1969 0::/services/uhttpd/instance1
3173 0::/services/snmpd/instance1
"""
# net-tools netstat on an old Debian: tcp6/udp6, and a program name with a space.
NET_TOOLS = """Active Internet connections (only servers)
Proto Recv-Q Send-Q Local Address Foreign Address State PID/Program name
tcp 0 0 0.0.0.0:22 0.0.0.0:* LISTEN 700/sshd: /usr/sbin
tcp6 0 0 :::22 :::* LISTEN 700/sshd: /usr/sbin
udp6 0 0 :::5353 :::* -
"""
def _netstat(out: str = BUSYBOX, cgroups: str = PROCD, *, rc: int = 0) -> str:
return f"SOCK_BEGIN\n[netstat]\n{out}__SS_RC={rc}\n[cgroups]\n{cgroups}SOCK_END\n"
class TestNetstat:
"""A host without ss (#673 in netOrk): OpenWrt's busybox, old net-tools."""
def test_a_socket_comes_with_its_process_and_procd_service(self):
sockets = _by_port(parse_listening_sockets(_netstat()))
assert sockets[("tcp", 22, "0.0.0.0")] == {
"proto": "tcp",
"address": "0.0.0.0",
"port": 22,
"interface": None,
"process": "dropbear",
"pid": 1604,
"unit": "dropbear",
"container_id": None,
}
@pytest.mark.parametrize(
"proto, port, address",
[("tcp", 22, "::"), ("tcp", 53, "fe80::1"), ("tcp", 53, "::1"), ("tcp", 53, "192.0.2.15")],
)
def test_every_address_form(self, proto, port, address):
assert (proto, port, address) in _by_port(parse_listening_sockets(_netstat()))
def test_a_udp_socket_has_no_state_column(self):
snmpd = _by_port(parse_listening_sockets(_netstat()))[("udp", 161, "0.0.0.0")]
assert (snmpd["process"], snmpd["pid"], snmpd["unit"]) == ("snmpd", 3173, "snmpd")
def test_a_socket_without_a_process_is_kept(self):
mdns = _by_port(parse_listening_sockets(_netstat()))[("udp", 5353, "0.0.0.0")]
assert (mdns["process"], mdns["pid"], mdns["unit"]) == (None, None, None)
def test_the_headers_are_no_sockets(self):
assert len(parse_listening_sockets(_netstat())) == 10
def test_net_tools_names_ipv6_and_programs_its_own_way(self):
sockets = _by_port(parse_listening_sockets(_netstat(NET_TOOLS, "")))
assert sockets[("tcp", 22, "::")]["process"] == "sshd"
assert sockets[("tcp", 22, "::")]["pid"] == 700
assert ("udp", 5353, "::") in sockets
def test_netstat_failing_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets(_netstat("netstat: invalid option -- 'p'\n", "", rc=1))
class TestTheNetstatFallback:
"""The command itself, on a host where ss is missing and netstat is not."""
@pytest.mark.skipif(
any(os.path.exists(f"{d}/ss") for d in ("/usr/sbin", "/sbin")),
reason="ss sits on the PATH the command always adds",
)
def test_it_reads_netstat_when_there_is_no_ss(self, tmp_path):
for tool in ("grep", "cut", "sort", "sed", "tr", "cat"):
(tmp_path / tool).symlink_to(shutil.which(tool))
stub = tmp_path / "netstat"
stub.write_text(f"#!/bin/sh\ncat <<'EOF'\n{BUSYBOX}EOF\n")
stub.chmod(0o755)
out = subprocess.run(
["/bin/sh", "-c", LISTENING_SOCKETS_COMMAND],
capture_output=True,
text=True,
env={"PATH": str(tmp_path)},
timeout=30,
).stdout
assert "[netstat]" in out
sockets = _by_port(parse_listening_sockets(out))
assert sockets[("tcp", 443, "0.0.0.0")]["process"] == "uhttpd"
assert len(sockets) == 10
class TestTheCommand:
def test_the_frame_is_not_in_the_command_itself(self):
"""An echoing transport prints the command back; the markers must only
appear once the command has run."""
assert "SOCK_BEGIN" not in LISTENING_SOCKETS_COMMAND
assert "SOCK_END" not in LISTENING_SOCKETS_COMMAND
def test_it_writes_and_changes_nothing(self):
for verb in ("sudo", " > ", ">>", "kill", "rm ", "systemctl "):
assert verb not in LISTENING_SOCKETS_COMMAND
def test_it_is_posix_sh(self):
result = subprocess.run(
["sh", "-n", "-c", LISTENING_SOCKETS_COMMAND], capture_output=True, text=True
)
assert result.returncode == 0, result.stderr
@pytest.mark.skipif(shutil.which("ss") is None, reason="needs ss on this host")
def test_it_runs_and_parses_on_this_host(self):
out = subprocess.run(
["sh", "-c", LISTENING_SOCKETS_COMMAND], capture_output=True, text=True, timeout=60
).stdout
sockets = parse_listening_sockets(out)
assert all(s["proto"] in ("tcp", "udp") and 0 < s["port"] < 65536 for s in sockets)
class _Driver(ListeningSocketsMixin):
def __init__(self, privileged: str, unprivileged: str = "") -> None:
self.answers = {True: privileged, False: unprivileged}
self.calls: list = []
def _run_listening_sockets_command(self, command: str, *, privileged: bool) -> str:
self.calls.append((command, privileged))
return self.answers[privileged]
class TestTheTemplate:
def test_a_driver_supplies_only_the_transport(self):
driver = _Driver(_wire())
reading = driver.get_listening_sockets()
assert reading["attributed"] is True
assert len(reading["sockets"]) == 9
[(command, privileged)] = driver.calls
assert privileged is True
def test_the_command_reaches_root_as_one_shell(self):
"""``sudo -n a; b`` runs only ``a`` as root: the whole script goes as
one ``sh -c`` argument."""
driver = _Driver(_wire())
driver.get_listening_sockets()
[(command, _privileged)] = driver.calls
assert command.startswith("sh -c '")
assert subprocess.run(["sh", "-n", "-c", command]).returncode == 0
def test_without_root_it_reads_what_it_can(self):
"""Without sudo, ss names only the user's own processes: the sockets
are still worth having, marked as not attributed."""
driver = _Driver("sudo: a password is required\n", _wire(cgroups=""))
reading = driver.get_listening_sockets()
assert reading["attributed"] is False
assert len(reading["sockets"]) == 9
assert [p for _c, p in driver.calls] == [True, False]
def test_a_host_without_either_is_not_asked_twice(self):
driver = _Driver("SOCK_BEGIN\n[no-ss]\nSOCK_END\n")
with pytest.raises(ListeningSocketsUnavailable):
driver.get_listening_sockets()
assert len(driver.calls) == 1
def test_not_every_os_driver_has_it(self):
"""A Windows host is an OSDriver too and has no ss: ``hasattr`` has to
stay a truthful answer, so the drivers that can mix this in themselves."""
assert not issubclass(OSDriver, ListeningSocketsMixin)
assert not hasattr(OSDriver, "get_listening_sockets")
+102
View File
@@ -0,0 +1,102 @@
"""What provisioning a VM from a cloud image needs to know about the guest,
the same for every hypervisor (NetOrk/netork#794)."""
from __future__ import annotations
import pytest
from napalm_device_types import HypervisorDriver
from napalm_device_types.provisioning import (
QEMU_GUEST_AGENTS,
network_config,
split_compression,
)
class TestGuestAgentDeclaration:
"""netOrk's cloud-init installs the agent through which the hypervisor
reads the new VM's IP. Which agent depends on the hypervisor and on the
guest; the guests a driver lists are the ones it can provision."""
def test_the_base_provisions_linux_with_qemu_guest_agent(self):
assert HypervisorDriver.GUEST_AGENTS == {
"linux": (("qemu-guest-agent",), ("systemctl enable --now qemu-guest-agent",)),
}
def test_the_declaration_is_data_netork_reads_off_the_class(self):
assert not callable(vars(HypervisorDriver)["GUEST_AGENTS"])
def test_the_per_hypervisor_attributes_are_gone(self):
assert not hasattr(HypervisorDriver, "GUEST_AGENT_PACKAGES")
assert not hasattr(HypervisorDriver, "GUEST_AGENT_RUNCMD")
class TestQemuGuestAgents:
"""Verified on FreeBSD 15.1 and OpenBSD 7.9 cloud images (#793)."""
def test_linux_is_the_base_default(self):
assert QEMU_GUEST_AGENTS["linux"] == HypervisorDriver.GUEST_AGENTS["linux"]
def test_freebsd(self):
assert QEMU_GUEST_AGENTS["freebsd"] == (
("qemu-guest-agent",),
("sysrc qemu_guest_agent_enable=YES", "service qemu-guest-agent start"),
)
def test_openbsd(self):
assert QEMU_GUEST_AGENTS["openbsd"] == (
("qemu-ga",),
("rcctl enable qemu_ga", "rcctl start qemu_ga"),
)
@pytest.mark.parametrize("guest_os", ["freebsd", "openbsd"])
def test_no_bsd_guest_is_told_to_use_systemd(self, guest_os):
_, runcmd = QEMU_GUEST_AGENTS[guest_os]
assert not any("systemctl" in command for command in runcmd)
class TestNetworkConfig:
"""cloud-init's network-config v2. FreeBSD's nuageinit reads only this
version; the v1 Proxmox generates makes it skip runcmd (#793)."""
def test_dhcp_nics_matched_by_mac(self):
cfg = network_config([("00:50:56:aa:bb:cc", True), ("00:50:56:aa:bb:dd", False)])
assert cfg == {
"version": 2,
"ethernets": {
"nic0": {"match": {"macaddress": "00:50:56:aa:bb:cc"}, "dhcp4": True},
},
}
def test_no_dhcp_nic_means_no_network_config(self):
assert network_config([("00:50:56:aa:bb:cc", False)]) is None
def test_macs_are_written_in_lower_case(self):
"""Proxmox reports them upper-case; nuageinit compares them as given."""
cfg = network_config([("BC:24:11:AA:BB:02", True)])
assert cfg["ethernets"]["nic0"]["match"]["macaddress"] == "bc:24:11:aa:bb:02"
class TestSplitCompression:
"""Official FreeBSD images come packed; hypervisors import unpacked ones."""
def test_xz(self):
assert split_compression("FreeBSD-15.1-RELEASE-amd64-BASIC-CLOUDINIT-ufs.qcow2.xz") == (
"FreeBSD-15.1-RELEASE-amd64-BASIC-CLOUDINIT-ufs.qcow2",
"xz -dc",
)
@pytest.mark.parametrize(
("name", "command"),
[("a.raw.gz", "gzip -dc"), ("a.raw.bz2", "bzip2 -dc"), ("a.qcow2.zst", "zstd -dc")],
)
def test_other_packers(self, name, command):
assert split_compression(name) == (name.rsplit(".", 1)[0], command)
@pytest.mark.parametrize(
"name", ["debian-13-genericcloud-amd64.qcow2", "ubuntu-24.04-server-cloudimg-amd64.img"]
)
def test_an_unpacked_image_is_left_as_it_is(self, name):
assert split_compression(name) == (name, None)
def test_the_suffix_is_matched_case_insensitively(self):
assert split_compression("IMAGE.QCOW2.XZ") == ("IMAGE.QCOW2", "xz -dc")
-17
View File
@@ -51,20 +51,3 @@ class TestVMConfigCarriesWhatAHardwareViewShows:
from napalm_device_types.models import VMPassthroughDict
assert get_type_hints(VMPassthroughDict) == {"slot": str, "kind": str, "config": str}
class TestGuestAgentDeclaration:
"""netOrk's cloud-init installed qemu-guest-agent on every new VM. A VMware
guest reports its IP through open-vm-tools instead; the hypervisor says
which, and netOrk stops hard-coding one of them."""
def test_default_is_qemu_guest_agent(self):
from napalm_device_types import HypervisorDriver
assert HypervisorDriver.GUEST_AGENT_PACKAGES == ("qemu-guest-agent",)
assert HypervisorDriver.GUEST_AGENT_RUNCMD == ("systemctl enable --now qemu-guest-agent",)
def test_attributes_are_not_methods(self):
from napalm_device_types import HypervisorDriver
assert not callable(vars(HypervisorDriver)["GUEST_AGENT_PACKAGES"])