feat: a public command channel and container engine access
CI / test (3.10) (push) Successful in 25s
CI / test (3.11) (push) Successful in 22s
CI / test (3.12) (push) Successful in 23s
CI / test (3.10) (pull_request) Successful in 24s
CI / test (3.11) (pull_request) Successful in 22s
CI / test (3.12) (pull_request) Successful in 22s
CI / test (3.10) (push) Successful in 25s
CI / test (3.11) (push) Successful in 22s
CI / test (3.12) (push) Successful in 23s
CI / test (3.10) (pull_request) Successful in 24s
CI / test (3.11) (pull_request) Successful in 22s
CI / test (3.12) (pull_request) Successful in 22s
netOrk reached a host's shell through napalm-linux's private `_send`: an interactive PTY with stdout and stderr merged and no exit code. For Docker it went further and opened its own paramiko connections around the driver. This adds the access layer only, no container model (NetOrk/netork#765): - `channel.py`: `CommandChannelMixin` declares `run_command()` (stdout, stderr, exit code) and `open_stream()` (a `ByteStream` to a running command). Declared under TYPE_CHECKING, so hasattr stays truthful. `run_on_transport()` and `open_stream_on_transport()` implement both on a paramiko exec channel for any SSH driver; `ParamikoExecStream` drains stderr on every read so the shared window never stalls. - `container_engine.py`: `ContainerEngineMixin` with `container_engines()` and `open_container_engine()`. The returned `ContainerEngineConnection` opens the engine API over `<binary> system dial-stdio` and runs the CLI with caller-chosen arguments. The driver decides the binary through `_container_engine_binary()`; callers never see the path. - README: why the container model lives in netOrk and not here. Additive; the Docker*Dict declarations stay until netOrk no longer reads them. Version 2.6.0. Refs #17
This commit is contained in:
@@ -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