Merge pull request 'feat: a public command channel and container engine access' (#18) from feat/container-engine-channel into main
This commit was merged in pull request #18.
This commit is contained in:
@@ -109,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.
|
||||
- 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
|
||||
|
||||
@@ -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`
|
||||
@@ -58,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
|
||||
@@ -104,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",
|
||||
|
||||
@@ -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,126 @@
|
||||
# -*- 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
|
||||
) -> CommandResult:
|
||||
"""Run the engine's CLI with *args* and wait for it."""
|
||||
return self._driver.run_command(self._command(args), timeout=timeout, stdin=stdin)
|
||||
|
||||
def stream_cli(self, args: Sequence[str], *, merge_stderr: 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``).
|
||||
"""
|
||||
command = self._command(args)
|
||||
return self._driver.open_stream(f"{command} 2>&1" if merge_stderr else command)
|
||||
|
||||
|
||||
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))
|
||||
@@ -871,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()``."""
|
||||
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "napalm-device-types"
|
||||
version = "2.5.0"
|
||||
version = "2.6.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,133 @@
|
||||
"""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")
|
||||
Reference in New Issue
Block a user