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")