From 735b683028d7cc4ef5d8beb4d56a57ac4fb68152 Mon Sep 17 00:00:00 2001 From: Christian Manivong Date: Wed, 7 Oct 2026 17:58:57 +0200 Subject: [PATCH] feat: a public command channel and container engine access netOrk reached a host's shell through napalm-linux's private `_send`: an interactive PTY with stdout and stderr merged and no exit code. For Docker it went further and opened its own paramiko connections around the driver. This adds the access layer only, no container model (NetOrk/netork#765): - `channel.py`: `CommandChannelMixin` declares `run_command()` (stdout, stderr, exit code) and `open_stream()` (a `ByteStream` to a running command). Declared under TYPE_CHECKING, so hasattr stays truthful. `run_on_transport()` and `open_stream_on_transport()` implement both on a paramiko exec channel for any SSH driver; `ParamikoExecStream` drains stderr on every read so the shared window never stalls. - `container_engine.py`: `ContainerEngineMixin` with `container_engines()` and `open_container_engine()`. The returned `ContainerEngineConnection` opens the engine API over ` system dial-stdio` and runs the CLI with caller-chosen arguments. The driver decides the binary through `_container_engine_binary()`; callers never see the path. - README: why the container model lives in netOrk and not here. Additive; the Docker*Dict declarations stay until netOrk no longer reads them. Version 2.6.0. Refs #17 --- README.md | 21 +++ napalm_device_types/__init__.py | 24 +++ napalm_device_types/base.py | 3 +- napalm_device_types/channel.py | 188 ++++++++++++++++++++++++ napalm_device_types/container_engine.py | 126 ++++++++++++++++ napalm_device_types/models.py | 11 ++ pyproject.toml | 2 +- tests/test_channel.py | 169 +++++++++++++++++++++ tests/test_container_engine.py | 133 +++++++++++++++++ 9 files changed, 675 insertions(+), 2 deletions(-) create mode 100644 napalm_device_types/channel.py create mode 100644 napalm_device_types/container_engine.py create mode 100644 tests/test_channel.py create mode 100644 tests/test_container_engine.py diff --git a/README.md b/README.md index 5ac73ee..337afc0 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/napalm_device_types/__init__.py b/napalm_device_types/__init__.py index fe92bd7..bb2bea3 100644 --- a/napalm_device_types/__init__.py +++ b/napalm_device_types/__init__.py @@ -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", diff --git a/napalm_device_types/base.py b/napalm_device_types/base.py index 110cefe..f6f8766 100644 --- a/napalm_device_types/base.py +++ b/napalm_device_types/base.py @@ -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 diff --git a/napalm_device_types/channel.py b/napalm_device_types/channel.py new file mode 100644 index 0000000..cb007ca --- /dev/null +++ b/napalm_device_types/channel.py @@ -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) diff --git a/napalm_device_types/container_engine.py b/napalm_device_types/container_engine.py new file mode 100644 index 0000000..b57035e --- /dev/null +++ b/napalm_device_types/container_engine.py @@ -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)) diff --git a/napalm_device_types/models.py b/napalm_device_types/models.py index 9ffa366..0b401ae 100644 --- a/napalm_device_types/models.py +++ b/napalm_device_types/models.py @@ -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()``.""" diff --git a/pyproject.toml b/pyproject.toml index a320c85..436e191 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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" diff --git a/tests/test_channel.py b/tests/test_channel.py new file mode 100644 index 0000000..0096b9b --- /dev/null +++ b/tests/test_channel.py @@ -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 diff --git a/tests/test_container_engine.py b/tests/test_container_engine.py new file mode 100644 index 0000000..ccaf73c --- /dev/null +++ b/tests/test_container_engine.py @@ -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")