Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
885c7e1f53 | ||
|
|
eb80d5cb0d | ||
|
|
e50e497939 | ||
|
|
3aed0b48d7 | ||
|
|
735b683028 | ||
|
|
08ec32e93a | ||
|
|
bd43bd75fa | ||
|
|
1112191aec | ||
|
|
b8b89acee1 | ||
|
|
f833e23422 | ||
|
|
b97ec654a0 | ||
|
|
c2d8d4a0d2 |
@@ -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/*
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
@@ -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,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
|
||||
|
||||
@@ -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)}
|
||||
@@ -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()``."""
|
||||
|
||||
|
||||
@@ -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
@@ -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"
|
||||
|
||||
@@ -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
|
||||
@@ -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)]
|
||||
@@ -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")
|
||||
@@ -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")
|
||||
@@ -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"])
|
||||
|
||||
Reference in New Issue
Block a user